Skip to main content

khive_db/stores/
graph.rs

1//! SQL-backed `GraphStore`: edge CRUD, neighbor queries, and bounded BFS traversal.
2
3#[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
36/// Map a rusqlite error to `StorageError` with `Graph` capability.
37fn 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
58// Latest eligible annotation of one target. First read at most three incident
59// `annotates` edge rows, tombstones included so the probe cannot scan a long
60// deleted prefix. The unique (namespace, source, target, relation) index allows
61// one row per source, so two or fewer rows is the complete incident set and
62// those notes are filtered and ordered directly. Three rows falls back to
63// walking notes newest-first, which keeps a long receipt history from being
64// enumerated or sorted. That fallback still visits newer unrelated notes when
65// a target with three or more incident edges has only old eligible annotations.
66const 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
98// The property predicate must run in both branches before ORDER BY/LIMIT;
99// filtering the single returned id in a caller lets a newer tagged decoy
100// conceal an older receipt with provenance.
101const 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
139// ---------------------------------------------------------------------------
140// Pure statement builders (ADR-099 B3 r6 structural cut) — see entity.rs's
141// sibling block for the full rationale. `upsert_edge`/`delete_edge` below and
142// `purge_incident_edges` (plus ADR-099's atomic prepare path in
143// `khive-runtime`, and `khive-runtime::operations::update_edge`'s
144// non-symmetric branch) all call these.
145// ---------------------------------------------------------------------------
146
147/// Replacement clause shared by canonical and guarded upserts. Live rows use
148/// replace semantics. Tombstones are excluded unless resurrection was
149/// explicitly requested; only that branch clears `deleted_at`.
150fn 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
170/// A `WHERE`-clause fragment asserting the id bound to `id_param` (an SQL
171/// placeholder like `?3`) resolves to a live edge endpoint — an
172/// undeleted entity or note, an event (append-only, no `deleted_at`), or an
173/// undeleted edge (the `annotates` relation's target may be any substrate,
174/// including another edge; ADR-002/ADR-055). Shared by
175/// [`edge_insert_guarded_by_endpoints_statement`] and
176/// [`edge_endpoints_exist`] (#769) so the two "does this endpoint still
177/// exist" probes — the guarded single-row insert and the guarded batch
178/// pre-check — can never drift on which substrates count as valid
179/// endpoints.
180fn 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
189/// Assert an exact live edge snapshot without updating its row or firing write
190/// triggers. Retained links additionally require their endpoints to remain live;
191/// deletion assertions may remove an annotation whose former target is gone.
192pub 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
218/// The exact natural-key-upserting `INSERT ... ON CONFLICT` this store's
219/// `upsert_edge` issues. Canonicalizes symmetric-relation endpoints first,
220/// matching `upsert_edge`'s own call to `canonical_edge_endpoints`.
221pub fn edge_upsert_statement(edge: &Edge) -> SqlStatement {
222    edge_upsert_statement_with_resurrection(edge, false)
223}
224
225/// Policy-aware form of [`edge_upsert_statement`].
226pub 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
274/// Insert a new edge only while both endpoints still exist.
275/// Competing IDs and natural keys, including tombstones, cause a constraint
276/// error rather than replacing or restoring the competing row. Callers must
277/// require one affected row to reject an endpoint removed after prepare.
278pub 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
293/// Conditional-insert companion to [`edge_upsert_statement`]. Conflicts on
294/// either the id or natural key leave the existing edge untouched so callers
295/// can read the winner and explicitly reapply their intended delta.
296pub 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
308/// Full-edge compare-and-swap update used after caller-side normalization
309/// was derived from a read snapshot. Unlike [`edge_upsert_statement`], this
310/// never inserts and cannot overwrite a row whose revision or deletion
311/// marker moved after the snapshot was read. The replacement revision must
312/// also be strictly greater than the persisted snapshot revision; equality
313/// is a refused CAS, never a successful write with an unchanged concurrency
314/// token. Mirrors `note_replace_if_unchanged_statement`
315/// (`crates/khive-db/src/stores/note.rs`). `created_at` is deliberately
316/// excluded from the `SET` list, matching the note/entity siblings — a
317/// replacement never rewrites the row's original creation time.
318pub 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/// The atomic `link` op's variant of [`edge_upsert_statement`] (ADR-099
367/// §B3). Shares the SAME
368/// `EDGE_NATURAL_KEY_CONFLICT_SET` conflict-arm text — the two builders
369/// cannot diverge on write behavior — but wraps the `INSERT` in a guarded
370/// `SELECT ... WHERE EXISTS(...)` that re-probes both endpoints for
371/// existence INSIDE the transaction, at commit time, rather than trusting
372/// prepare-time validation alone.
373///
374/// This guard is atomic-`link`-specific, not a `edge_upsert_statement`
375/// concern: `LinkPlan`'s own doc comment (`khive-runtime::atomic_plan`)
376/// records why it must be commit-time, not prepare-time — a `link` op's
377/// async prepare pass (`validate_edge_relation_endpoints`) can run and pass
378/// BEFORE an earlier op in the SAME atomic unit (e.g. `delete(X, hard)`)
379/// removes that very endpoint; only a commit-time, in-transaction guard
380/// closes that intra-batch ordering hazard (ADR-099 acceptance criteria:
381/// `[delete(X, hard), link(A, X)]` must fail, not silently create a
382/// dangling edge). Canonical `link` has no equivalent need — it executes
383/// and commits standalone, with no other op's write interleaved between its
384/// own validation and its own write.
385#[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/// Policy-aware form of [`edge_insert_guarded_by_endpoints_statement`].
402#[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/// Atomic `link` create arm. Unlike the compatibility upsert builder above,
445/// this statement never changes a row that appeared after prepare: a natural-
446/// key conflict affects zero rows and therefore fails the caller's guard.
447#[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
485/// Atomic `link` replace/resurrection arm. The natural-key row and both
486/// concurrency markers must still match the prepare snapshot, and both
487/// endpoints must still exist. Clearing `deleted_at` is deliberate here:
488/// callers only build this statement for a live replacement or an explicitly
489/// authorized resurrection.
490pub 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
529/// The exact soft-delete `UPDATE` this store's `delete_edge(Soft)` issues.
530pub 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
540/// The exact hard-delete `DELETE` this store's `delete_edge(Hard)` issues.
541pub 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
549/// The exact cascade `DELETE` this store's `purge_incident_edges` issues.
550pub 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
558// ---------------------------------------------------------------------------
559// Symmetric-relation update DML (ADR-099 B3 r6 second pass) — the SQL text
560// `khive-runtime::operations::KhiveRuntime::update_edge_symmetric_dml` (the
561// synchronous raw-connection commit-time path, run inside the writer-task/
562// pool-mutex transaction) and ADR-099's atomic `prepare_update_edge` symmetric
563// branch (the async plan-time path) both bind. `upsert_edge` cannot be used
564// here: it resolves `ON CONFLICT(namespace, id)` first and cannot detect a
565// natural-key collision at (namespace, source_id, target_id, relation) with a
566// *different* id, which is exactly the case a symmetric-relation endpoint
567// canonicalization can produce.
568//
569// The two call sites bind these against different parameter-passing
570// mechanisms — `conn.execute`/`conn.query_row` with `rusqlite::params!` in the
571// synchronous path (it must run inside an existing transaction on a borrowed
572// `&rusqlite::Connection`, so it cannot go through the `SqlStatement`/
573// `SqlValue` plan-shape khive-storage abstracts elsewhere) vs. `SqlValue`
574// plan params for the async `PlanStatement` path — but the SQL TEXT itself
575// (the `EDGE_SYMMETRIC_*_SQL` constants below) is the single source of truth
576// for both, closing the class of drift that produced a hand-copied SQL
577// literal silently diverging from canonical (ADR-099 §B3).
578pub 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
585/// Canonical `update_edge`'s guarded variant of
586/// [`EDGE_SYMMETRIC_DELETE_NONCANONICAL_SQL`]: `?3`/`?4` pin the fetched
587/// snapshot's `updated_at`/`deleted_at` so a writer whose edge changed
588/// concurrently after it read that snapshot cannot delete the row out from
589/// under the concurrent write, even though a canonical survivor genuinely
590/// exists at the natural key. Zero affected rows means stale, not
591/// "no conflict" — the caller has already confirmed a conflicting canonical
592/// row exists before running this statement. Deliberately a DIFFERENT
593/// constant from the unguarded one above: merge's predicate-based rewrites
594/// (`khive-runtime::curation`) intentionally keep running the unguarded form
595/// inside their own single writer transaction and must not be changed to
596/// bind this one.
597pub 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
601/// Case (a) update, guarded on the fetched snapshot's revision and deletion
602/// marker: `?9`/`?10` pin `updated_at`/`deleted_at` as read, and `?5 >
603/// updated_at` requires the replacement revision to strictly advance —
604/// mirroring `edge_replace_if_unchanged_statement`'s guard so the symmetric
605/// path cannot silently overwrite a concurrent writer's change between the
606/// snapshot read and this write.
607pub 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
614/// Plan-shape builder for [`EDGE_SYMMETRIC_CONFLICT_PROBE_SQL`] — the
615/// async prepare-time conflict probe.
616pub 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
636/// Plan-shape builder for [`EDGE_SYMMETRIC_DELETE_NONCANONICAL_SQL`] —
637/// case (b): a canonical row already exists, delete the requested row.
638pub 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/// Plan-shape builder for [`EDGE_SYMMETRIC_UPDATE_INPLACE_SQL`] —
650/// case (a): no conflict, update the requested row in place, guarded on the
651/// fetched snapshot's revision and deletion marker.
652#[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// ---------------------------------------------------------------------------
690// Symmetric-relation update DML — atomic-only, commit-time self-guarding
691// variant (ADR-099 §B3).
692//
693// Canonical `update_edge_symmetric_dml` binds the shared SQL constants above
694// directly, not the plan-shape builders — the builders are used by the atomic
695// prepare path. Canonical probes and branches synchronously INSIDE its own
696// writer-task transaction, with no other op interleaved between its probe and
697// its write, so the interleaving exposure the atomic path has does not arise
698// there. Canonical's absorption delete is nonetheless bound to the GUARDED
699// constant, because its snapshot is read before the transaction and can be
700// stale by the time the delete runs.
701//
702// The atomic path is structurally different: its conflict probe runs in the
703// async PREPARE phase, which for a multi-op `--atomic` unit completes for
704// EVERY op before the synchronous COMMIT phase begins for ANY of them. An
705// earlier op in the SAME atomic unit (e.g. a `delete` or another symmetric
706// `update` touching the same natural key) can change the conflict landscape
707// between this op's prepare-time probe and its own statements finally
708// executing at commit time — this is a real staleness
709// window, not just an SQL-text duplication concern.
710//
711// The two builders below close it: instead of a Rust-level `if let
712// Some(conflict) { ... } else { ... }` that hand-picks ONE of three
713// statements at prepare time (the second hand-assembled branch this
714// closes), the atomic plan ALWAYS carries both statements, in order, and
715// each is a self-guarding, commit-time predicate that re-evaluates conflict
716// state fresh against whatever the transaction's state actually is when it
717// runs — not what prepare's probe said:
718//
719// 1. [`edge_symmetric_delete_if_conflict_statement`]: deletes the requested
720//    (non-canonical) row IF AND ONLY IF a differently-id'd canonical row
721//    exists at the target natural key at THIS moment AND the row's own
722//    `updated_at`/`deleted_at` still match the snapshot this plan was built
723//    from (guard: 0 or 1 rows) — a plan built from a since-changed snapshot
724//    must not delete the row just because some other, unrelated conflict
725//    happens to exist; the second statement's own commit-time predicate
726//    (below) then sees `changes() = 0` and fails its `exactly(1)` guard,
727//    aborting the whole atomic unit rather than silently absorbing a
728//    concurrent writer's change.
729// 2. [`edge_symmetric_absorb_or_update_inplace_statement`]: a single
730//    `UPDATE` that no longer trusts an `id = ?2 OR natural-key` predicate
731//    (ADR-099 §B3 — that predicate could
732//    match the WRONG row: if a different op earlier in the SAME atomic unit
733//    had already deleted the requested edge, statement 1 above no-ops (its
734//    "0 rows" result is indistinguishable at the Rust level from "no
735//    conflict existed"), yet the natural-key arm could still hit a
736//    pre-existing canonical row that this update never causally touched).
737//    The fix ties the natural-key arm to `changes()` — SQLite's per-
738//    connection scalar reporting the row count of the most recently
739//    COMPLETED statement, i.e. statement 1's own result, evaluated fresh at
740//    THIS statement's execution, not at prepare time:
741//    - `id = ?2 AND changes() = 0`: statement 1 deleted nothing, so the
742//      requested row is still live under its own id — update it in place.
743//    - `source_id = ?3 AND target_id = ?4 AND relation = ?5 AND id != ?2
744//      AND changes() = 1`: statement 1 just deleted the requested row
745//      BECAUSE a conflict existed — match the surviving canonical row and
746//      leave its attributes unchanged (ADR-039 DO NOTHING; see below).
747//    These two arms are mutually exclusive and, together with statement 1's
748//    own guard, jointly exhaustive: if the requested row no longer existed
749//    when statement 1 ran (the same-unit race above), statement 1 affects 0
750//    rows for a reason unrelated to conflict absorption, `id = ?2` no
751//    longer matches anything (the row is gone), and the natural-key arm's
752//    `changes() = 1` guard is false — so this statement affects ZERO rows
753//    and the plan's `AffectedRowGuard::exactly(1)` on it fails the op,
754//    aborting the whole atomic unit rather than silently mutating an
755//    unrelated row.
756//
757//    ADR-039's edge-conflict contract is ON CONFLICT DO NOTHING: the
758//    natural-key (absorbed-conflict) arm must leave the surviving canonical
759//    row's attributes exactly as they were — refreshing them from the
760//    discarded edge (and forcing `deleted_at = NULL`, resurrecting a
761//    tombstone) is the same defect already fixed on the merge-rewire path
762//    (khive#1213). Every SET expression is therefore keyed on `id = ?2`:
763//    that condition is true only for the WHERE clause's first (in-place)
764//    arm — the second (absorbed) arm only ever matches a row whose id is
765//    NOT ?2 — so each column either takes its new value (in-place arm) or
766//    self-assigns its current value (absorbed arm, a true no-op). This
767//    still affects exactly one row either way (SQLite's `changes()` counts
768//    matched rows, not changed bytes), so `AffectedRowGuard::exactly(1)`
769//    and the race-abort behavior above are unaffected by this being a
770//    no-op write in the absorbed case.
771//
772// No probe, no branch, no read at all is needed to APPLY this pair. Which
773// row this plan actually touched is derived post-commit by the caller via a
774// fresh natural-key lookup (`khive-runtime::KhiveRuntime::get_edge_by_natural_key_including_deleted`,
775// filtered on the canonicalized endpoints/relation and including soft-deleted rows — unlike
776// `list_edges`, which unconditionally filters `deleted_at IS NULL` and would report "not
777// found" for a surviving row this absorption arm left tombstoned) — ADR-099 §B3
778// removed the prior prepare-time advisory `target_id` probe entirely:
779// a value computed before the
780// SAME atomic unit's other ops have run is not a fact this plan can stand
781// behind, so result rendering no longer trusts it.
782#[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/// `?10`/`?11` pin the fetched snapshot's `updated_at`/`deleted_at` and
819/// `?7 > updated_at` requires the replacement revision to strictly advance —
820/// guarding the in-place arm (`id = ?2`) against a concurrent writer that
821/// changed the row between the snapshot read and this statement. The
822/// absorbed arm (a differently-id'd canonical row) is intentionally left
823/// unguarded on revision: per ADR-039 DO NOTHING it self-assigns every
824/// column, so it is a no-op write regardless of the survivor's current
825/// state.
826#[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/// Internal seam for khive-runtime; no compatibility promise.
885///
886/// A live source-document reference observed before pack preparation.
887#[doc(hidden)]
888pub struct GraphDocumentGuard {
889    pub namespace: String,
890    pub id: Uuid,
891    pub expected_blob_ref: String,
892}
893
894/// Internal seam for khive-runtime; no compatibility promise.
895///
896/// The complete row, or absence, observed for one canonical natural key.
897#[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/// Internal seam for khive-runtime; no compatibility promise.
907///
908/// Source-row expectations checked before any graph mutation.
909#[doc(hidden)]
910#[derive(Default)]
911pub struct GraphMutationPreconditions {
912    pub document: Option<GraphDocumentGuard>,
913    pub edges: Vec<GraphEdgeSnapshotGuard>,
914}
915
916/// Internal seam for khive-runtime; no compatibility promise.
917///
918/// Existing graph engines selected for a source-backend composition unit.
919#[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/// Internal seam for khive-runtime; no compatibility promise.
936///
937/// Transaction-observed outcomes retain the existing graph classifications.
938#[doc(hidden)]
939pub enum GraphMutationOutcome {
940    Single(GuardedEdgeUpsertOutcome),
941    Batch(GuardedEdgeBatchUpsertOutcome),
942    CommitAnnotation(CommitAnnotationInsertOutcome),
943}
944
945/// Internal seam for khive-runtime; no compatibility promise.
946///
947/// Ordered upsert results and the separate preimages of retired live edges.
948#[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/// Internal seam for khive-runtime; no compatibility promise.
970///
971/// Guarded composition seam for packs.
972#[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    // Retain the normal accessors' store-readiness checks before admission;
990    // no schema work or asynchronous dispatch enters the enlisted engine.
991    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
1135/// Connection-enlisted graph/event engine: no transaction control or await.
1136fn 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
1239/// A GraphStore backed by SQLite tables.
1240pub struct SqlGraphStore {
1241    pool: Arc<ConnectionPool>,
1242    index_repair: Option<super::index_repair::IndexRepairContext>,
1243    is_file_backed: bool,
1244    /// Default namespace for multi-record queries (ADR-007 PARAM-ONLY: used as a
1245    /// WHERE filter on `query_edges`/`neighbors`/`traverse`, never as an
1246    /// enforcement gate on by-ID operations).
1247    namespace: String,
1248    writer_task: Option<WriterTaskHandle>,
1249}
1250
1251impl SqlGraphStore {
1252    /// Create a new store with a default namespace for multi-record query filtering.
1253    ///
1254    /// The namespace is a PARAM-ONLY hint (ADR-007 rule 4) — it is used as a
1255    /// WHERE filter in multi-record queries and as the write namespace stamped on
1256    /// upserted edges, but it does NOT enforce isolation: `upsert_edge` accepts
1257    /// edges from any namespace, and by-ID ops (`get_edge`, `delete_edge`) ignore
1258    /// the namespace entirely.
1259    pub fn new_scoped(
1260        pool: Arc<ConnectionPool>,
1261        is_file_backed: bool,
1262        namespace: impl Into<String>,
1263    ) -> Self {
1264        // Enabled by default for file-backed pools; explicit off/degraded
1265        // construction remains synchronous (ADR-067 Component A, mirrors
1266        // entity.rs policy): a missing writer task is cached without failing
1267        // construction. Every write re-resolves it and applies
1268        // strict/compatibility policy then.
1269        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    /// Route a single-row write through the pool-wide `WriterTask` when
1312    /// the write queue is enabled and a handle is available. Strict mode
1313    /// refuses a missing handle; compatibility mode falls back to the legacy
1314    /// standalone-connection / pool-mutex path (ADR-067 Component A, Fork C
1315    /// slice 2). `f` must be DML-only. See
1316    /// `crates/khive-db/docs/api/graph.md` for the per-caller routing rules.
1317    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
1410// =============================================================================
1411// Helpers
1412// =============================================================================
1413
1414/// Publish a graph read's accounting after the read resolves, on either
1415/// outcome.
1416///
1417/// Both counters are collected inside the storage closure because that closure
1418/// runs on a blocking thread where the task-local usage context is invisible,
1419/// and because a value returned only on success reports nothing at all when
1420/// the read fails partway. `queries` is incremented immediately before a
1421/// statement executes, so a read that never got that far — a connection the
1422/// pool could not hand out, a blocking task that died — counts no round trip.
1423/// `rows` is incremented as each adjacency entry materialises, so entries
1424/// storage already returned still count when a later row fails to convert.
1425fn 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
1491/// Build the `relation IN (...)` / `weight >= ?` `WHERE`-extra clause and the
1492/// `LIMIT` clause shared by `neighbors` and `neighbors_both_directions` —
1493/// both filter and cap identically, differing only in which direction(s) the
1494/// base `SELECT`s cover. `start_param_idx` is the next free `?N` placeholder
1495/// (both callers bind `namespace` and `node_id` as `?1`/`?2` first).
1496fn 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// Test-only counter of storage-level neighbor SELECT executions (`neighbors`
1585// and `neighbors_both_directions` each issue exactly one `graph_edges`
1586// query per call). Lets tests assert the query-count halving a
1587// `Direction::Both` caller gets from `neighbors_both_directions` vs the old
1588// pattern of two separate `neighbors` calls (ADR-089 context-verb
1589// optimization). Gated out of release builds — no counter overhead on the
1590// hot path in production.
1591#[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
1619/// Deterministic ORDER BY for edge queries. #1671: an `id` tiebreak is
1620/// always appended (following the last sort field's direction, DESC for the
1621/// empty default) so pages over equal sort values keep a total order. That
1622/// removes tie-order instability only — offset paging can still duplicate or
1623/// skip rows under concurrent inserts/deletes or sort-key updates (that
1624/// would need snapshot isolation or keyset pagination).
1625fn 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
1647/// `CASE` expression classifying one edge endpoint by the base it resolves
1648/// against. Soft-deleted endpoints are already excluded by
1649/// [`LIVE_ENDPOINTS_CONDITION`] in the surrounding `WHERE`, so a row reaching
1650/// this expression has live endpoints or none at all.
1651///
1652/// The probe reads the same local tables as [`LIVE_ENDPOINTS_CONDITION`] and
1653/// [`endpoint_exists_clause`], which is what makes the three agree on what an
1654/// endpoint is. Two consequences are deliberate. An `annotates` endpoint that
1655/// is an event or another edge is a live endpoint that is not in either base,
1656/// so it classifies as `none` and lands in `unresolved`; that is the bucket's
1657/// meaning, not a miscount, and `annotates` is excluded from structural
1658/// density anyway. And when ADR-009 edge routing is wired, an endpoint held in
1659/// another backend will not be in these tables either: `target_backend` is
1660/// written by nothing today, so no such row exists yet, but whoever wires the
1661/// routing has to decide what the breakdown should say about it and change
1662/// this expression deliberately rather than discover it as a silent
1663/// `unresolved`.
1664fn 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
1673/// Folds one `(source_base, target_base, count)` group into the tally. An
1674/// unrecognized pair lands in `unresolved` rather than being dropped, so the
1675/// buckets keep summing to the live edge total.
1676fn 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
1687/// Restricts an edge count to edges whose endpoints are still live.
1688///
1689/// A soft delete leaves incident edges in place on purpose (`KhiveRuntime::delete_entity`
1690/// documents it; only a hard delete purges them). Every reader that walks the graph
1691/// hydrates the endpoint records, so it cannot reach an edge whose endpoint is
1692/// tombstoned. A count that includes those edges therefore reports a density no reader
1693/// can walk, under a `count_scope` that says `live_only`.
1694///
1695/// The predicate excludes an endpoint only when it is PRESENT AND tombstoned. An id
1696/// found in neither table belongs to a substrate this database does not hold, and
1697/// treating its absence as a tombstone would silently under-count; erring toward
1698/// counting keeps the failure in the direction the old behaviour already had.
1699const 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
1708/// Append [`LIVE_ENDPOINTS_CONDITION`] to a `WHERE` clause produced by the edge filter
1709/// builders, which always emit a non-empty ` WHERE ...`.
1710fn 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
1741/// Same filter conditions as [`build_edge_filter_sql_for_namespaces`], but
1742/// the namespace set is bound as a single pre-serialized JSON array
1743/// parameter (`?1`) matched via `json_each` rather than one `?N` per
1744/// namespace.
1745///
1746/// `namespace IN (?1, ?2, ..., ?N)` costs one bound SQLite variable per
1747/// namespace; a caller with a large visible-namespace set (hundreds to
1748/// thousands, e.g. a broad `[actor]` visibility grant) can exceed
1749/// `SQLITE_LIMIT_VARIABLE_NUMBER` (999 by default) before any of the
1750/// filter's own parameters are even added. Binding one JSON string instead
1751/// keeps this at O(1) parameters regardless of namespace count, so paged
1752/// enumeration (which needs one exact-order SQL statement per page — see
1753/// [`SqlGraphStore::query_edges_in_namespaces`]) never has to chunk the
1754/// namespace set and cannot silently corrupt paging by doing so (#2088).
1755fn 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
1852// =============================================================================
1853// GraphStore implementation
1854// =============================================================================
1855
1856/// Canonical endpoint order for symmetric relations (F012).
1857///
1858/// For `competes_with` and `composed_with`, ensures `source_uuid < target_uuid`
1859/// so A→B and B→A collapse to a single canonical row in storage.
1860fn 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
1868/// Standalone existence probe for both endpoints of a would-be edge (#769),
1869/// matching the `WHERE EXISTS(...)` shape
1870/// [`edge_insert_guarded_by_endpoints_statement`] embeds in its own guarded
1871/// `INSERT`. Returns per-endpoint existence (not a single AND'd bool) so
1872/// callers can report exactly which side was missing. See
1873/// `crates/khive-db/docs/api/graph.md` for its two call sites and why both
1874/// need in-transaction, not reconstructed-after-the-fact, results.
1875fn 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
1934/// The predicate and insert are one DML statement inside the caller's writer
1935/// transaction. The follow-up probes classify an unsuccessful insert under
1936/// that same transaction, so no racing curation can turn a tombstone into a
1937/// repairable gap or change the note identity after its check.
1938fn conditional_commit_annotation_insert(
1939    conn: &rusqlite::Connection,
1940    edge: Edge,
1941    guard: &CommitAnnotationGuard,
1942) -> Result<CommitAnnotationInsertOutcome, rusqlite::Error> {
1943    // This connection is inside BEGIN IMMEDIATE on both writer paths. Check
1944    // the pack-owned cursor rows here, then bind the result into the core-only
1945    // INSERT: no other writer can change them between this probe and the DML.
1946    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
2028/// Apply one replacement-style upsert and derive its disposition/preimage on
2029/// the same write connection. Callers keep this inside one write transaction.
2030fn 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        // Test-only rendezvous proving the endpoint decision is made while
2042        // this write transaction still excludes racing writers.
2043        #[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
2088/// Shared DML-only batch engine for observed and legacy upserts on both writer
2089/// routes. The caller owns the transaction and rolls back SQL failures.
2090fn observed_edge_batch_upsert(
2091    conn: &rusqlite::Connection,
2092    requests: &[EdgeUpsertRequest],
2093    guard_endpoints: bool,
2094) -> Result<GuardedEdgeBatchUpsertOutcome, rusqlite::Error> {
2095    // Preflight every refusal before the first mutation so the successful
2096    // return path is genuinely all-or-nothing even under WriterTask, whose
2097    // transaction commits any `Ok(...)` closure result.
2098    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
2144/// One index-seek adjacency statement used by bounded level-synchronous BFS.
2145/// There is deliberately no `ORDER BY`: the public contract promises minimum
2146/// depth, not a same-depth tie order, and sorting a high-degree node before a
2147/// low `LIMIT` would reintroduce the unbounded SQL work fixed by #1444.
2148fn 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                        // Test seam after `sqlite3_step` has produced a row:
2328                        // the cursor's read snapshot is demonstrably live.
2329                        #[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                    // Keep the culprit's legacy first_error while recording
2720                    // per-input error classes and retryability in input order.
2721                    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        // SQLite SQLITE_MAX_VARIABLE_NUMBER defaults to 999; chunk at 900 to stay safe.
2869        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        // Declared across the whole chunk loop, not per chunk: a failure in a
2973        // later chunk must not erase what the earlier chunks already did.
2974        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                            // min_weight is returned separately so it can be added to
3023                            // all_params AFTER the rel_params block, at the right index.
3024                            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                    // Wrap combined inner with per-source ROW_NUMBER limit if needed.
3076                    //
3077                    // Deterministic weight-descending order, tie-broken by node_id
3078                    // ascending, applied INSIDE the window's ORDER BY — otherwise a
3079                    // per-origin cap can silently drop high-weight neighbors in favor
3080                    // of arbitrary SQLite row order (mirrors neighbors(), ADR-089
3081                    // context-verb review; issue #589).
3082                    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        // One statement with `namespace IN (...)` and real SQL paging: a
3305        // per-namespace prefix fetch merged and re-sliced client-side floats
3306        // the offset window between calls, silently duplicating and skipping
3307        // rows during enumeration (#2088).
3308        //
3309        // The namespace set is bound as a single JSON-array parameter (see
3310        // `build_edge_filter_sql_for_namespaces_json`), not one `?N` per
3311        // namespace: a caller visible in hundreds-to-thousands of namespaces
3312        // would otherwise blow past `SQLITE_LIMIT_VARIABLE_NUMBER` before any
3313        // filter parameter is even added, and chunking the namespace set (as
3314        // the count-only aggregate methods below do) is not an option here —
3315        // it would reintroduce exactly the floating-offset bug this method
3316        // exists to close, since a per-chunk LIMIT/OFFSET cannot be
3317        // re-sliced into one globally exact page.
3318        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            // `where_clause` always pins `namespace = ?1`; adding `id > ?N` here
3533            // keeps the predicate a range scan against the implicit unique index
3534            // backing `PRIMARY KEY (namespace, id)` — equality on the leading
3535            // column plus a range on the trailing one, with `ORDER BY id ASC`
3536            // matching the index order, so SQLite seeks instead of scanning.
3537            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            // CROSS JOIN fixes the ledger as the outer loop, preserving an
3589            // indexed `seq > boundary` scan with no full-match sort.
3590            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    /// Single-query both-direction neighbor fetch (ADR-089 context-verb
3644    /// optimization): projects a `'out'`/`'in'` literal from each `UNION ALL`
3645    /// arm so the caller gets direction labels without a second direction-
3646    /// scoped round trip. `query.direction` is ignored — always both.
3647    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                // Same global weight-descending/node_id-ascending order as `neighbors`
3677                // (ADR-089 context-verb review),
3678                // applied across BOTH directions before `LIMIT` truncates. A
3679                // reciprocal pair (an Out edge and an In edge to/from the same
3680                // neighbor at the same weight) ties on `(weight, node_id)`, so the
3681                // order is extended with a direction rank (`out` before `in`) and
3682                // finally `edge_id` to make the pre-`LIMIT` order fully
3683                // deterministic).
3684                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                // Prepare and step the actual indexed adjacency statements with
3788                // LIMIT 0 before BFS takes its long-lived reader. This checks the
3789                // schema cookie without consuming walk rows or query counters.
3790                // A schema change during the walk remains its original typed
3791                // failure; we never restart BFS or renew its budget.
3792                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                // Dropping preflight stops pending repair admission. Already admitted
3805                // constructor DDL keeps its completion ownership outside this await.
3806                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            // Shared with the blocking closure so the counts survive an error.
3843            // `with_reader` runs the closure on a blocking thread where the
3844            // task-local usage context is invisible, so it cannot call
3845            // `usage::count` itself; returning the totals in the Ok value would
3846            // lose every round trip already issued when a later statement fails.
3847            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            // Accounted on BOTH outcomes. `db_round_trips` counts round trips
3877            // *issued* and `graph_hops` counts adjacency rows storage *returned*,
3878            // so work already done before a later statement errors is real work and
3879            // must appear. Reading the shared counters here rather than off the
3880            // Ok value is what makes that possible: the closure runs on a blocking
3881            // thread where the task-local usage context is not visible, so it
3882            // cannot count for itself, and a value returned only on success
3883            // reports nothing at all when the traversal fails partway.
3884            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        // No namespace filter: UUID v4 is globally unique. Hard-delete cascade must
3900        // remove ALL incident edges regardless of which namespace they were written in
3901        // (ADR-002 no-dangling-references, ADR-007 by-ID contract).
3902        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
3912// =============================================================================
3913// DDL
3914// =============================================================================
3915
3916const 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;