Skip to main content

khive_storage/
graph.rs

1//! Graph storage capability — edge CRUD and traversal.
2
3use async_trait::async_trait;
4use chrono::{DateTime, Utc};
5use khive_types::EdgeRelation;
6use uuid::Uuid;
7
8use crate::capability::StorageCapability;
9use crate::error::StorageError;
10use crate::types::{
11    BatchWriteSummary, DeleteMode, DirectedNeighborHit, Direction, Edge, EdgeEndpointBaseCounts,
12    EdgeFilter, EdgeSeekPage, EdgeSortField, EdgeUpsertRequest, EdgeUpsertResult, GraphPath,
13    GuardedBatchOutcome, GuardedEdgeBatchUpsertOutcome, GuardedEdgeUpsertOutcome,
14    GuardedWriteOutcome, LinkId, NeighborCursor, NeighborHit, NeighborQuery, Page, PageRequest,
15    SeekCursor, SeekPage, SortOrder, StorageResult, TraversalRequest,
16};
17
18/// Directed edge CRUD and graph traversal over the knowledge graph.
19#[async_trait]
20pub trait GraphStore: Send + Sync + 'static {
21    /// Return the newest live note of `kind` carrying the exact string `tag`
22    /// in its properties.tags array and connected to `node_id` by a live
23    /// incoming `annotates` edge in this store's namespace. Returns its UUID
24    /// and creation timestamp; equal timestamps choose the smallest UUID.
25    ///
26    /// Apply all predicates before limiting to one result. Note lookup follows
27    /// the by-ID contract; visibility is determined by the annotation edge's
28    /// namespace, not by introducing a second namespace filter on the note.
29    /// The note and edge must belong to this backend. Unsupported backends
30    /// fail explicitly rather than scanning an arbitrary annotation window.
31    async fn latest_annotating_note(
32        &self,
33        _node_id: Uuid,
34        _kind: &str,
35        _tag: &str,
36    ) -> StorageResult<Option<(Uuid, i64)>> {
37        Err(StorageError::Unsupported {
38            capability: StorageCapability::Graph,
39            operation: "latest_annotating_note".into(),
40            message: "this backend does not implement latest matching annotation lookup".into(),
41        })
42    }
43
44    /// Insert or update a single edge.
45    async fn upsert_edge(&self, edge: Edge) -> StorageResult<()>;
46    /// Insert an edge only when neither its id nor natural key already
47    /// exists. Returns `true` when this call inserted the row and `false`
48    /// when an existing row won the race. The existing row is never updated.
49    ///
50    /// The default returns `Unsupported` rather than falling back to
51    /// [`GraphStore::upsert_edge`], because an upsert would overwrite the
52    /// winning row and violate this method's conditional-insert contract.
53    async fn insert_edge_if_absent(&self, _edge: Edge) -> StorageResult<bool> {
54        Err(StorageError::Unsupported {
55            capability: StorageCapability::Graph,
56            operation: "insert_edge_if_absent".into(),
57            message: "this backend does not implement conditional edge insert".into(),
58        })
59    }
60    /// Insert or update a batch of edges.
61    async fn upsert_edges(&self, edges: Vec<Edge>) -> StorageResult<BatchWriteSummary>;
62    /// Insert or replace one edge and return the transaction-observed
63    /// disposition plus preimage. Tombstone restoration is controlled by the
64    /// request rather than being an implicit side effect of every upsert.
65    async fn upsert_edge_observed(
66        &self,
67        _request: EdgeUpsertRequest,
68    ) -> StorageResult<EdgeUpsertResult> {
69        Err(StorageError::Unsupported {
70            capability: StorageCapability::Graph,
71            operation: "upsert_edge_observed".into(),
72            message: "this backend does not implement observed edge upserts".into(),
73        })
74    }
75    /// Replace an edge only when the persisted row still matches the
76    /// caller's read snapshot.
77    ///
78    /// `expected_updated_at` is the snapshot revision and
79    /// `expected_deleted_at` closes the soft-delete race. The replacement
80    /// edge's `updated_at` must be strictly greater than that persisted
81    /// revision. Returns `false` when the row disappeared, changed, or was
82    /// supplied a non-advancing replacement revision. This is the full-edge
83    /// compare-and-swap seam used when a caller derives coupled fields from
84    /// that snapshot before persistence — mirrors
85    /// [`crate::NoteStore::replace_note_if_unchanged`]. The default returns
86    /// `Unsupported` rather than falling back to an unguarded upsert and
87    /// reintroducing the stale-snapshot race.
88    async fn replace_edge_if_unchanged(
89        &self,
90        _edge: Edge,
91        _expected_updated_at: DateTime<Utc>,
92        _expected_deleted_at: Option<DateTime<Utc>>,
93    ) -> StorageResult<bool> {
94        Err(StorageError::Unsupported {
95            capability: StorageCapability::Graph,
96            operation: "replace_edge_if_unchanged".into(),
97            message: "this backend does not implement guarded edge replacement".into(),
98        })
99    }
100    /// Insert or update a single edge, re-checking that both endpoints still
101    /// exist (and are not soft-deleted) as part of the same write, not a
102    /// separate prior read. Closes the TOCTOU window between an async
103    /// prepare-time existence check and a later, unconditional write: a
104    /// concurrent hard-delete of an endpoint that lands between the two can
105    /// otherwise leave a durably dangling edge (#769).
106    ///
107    /// Returns [`GuardedWriteOutcome::Refused`] naming exactly which
108    /// endpoint(s) were missing, determined by the guard's own in-transaction
109    /// probe — never reconstructed by a caller re-reading the endpoints after
110    /// the write already failed, since a concurrent write landing between the
111    /// refusal and any such later read could misreport which endpoint was
112    /// actually missing at write time.
113    ///
114    /// Default returns `StorageError::Unsupported`: a backend that does not
115    /// override this method cannot honor the endpoint-existence guarantee,
116    /// and silently falling back to [`GraphStore::upsert_edge`] would
117    /// reintroduce the TOCTOU window this method exists to close.
118    async fn upsert_edge_guarded(&self, _edge: Edge) -> StorageResult<GuardedWriteOutcome> {
119        Err(StorageError::Unsupported {
120            capability: StorageCapability::Graph,
121            operation: "upsert_edge_guarded".into(),
122            message: "this backend does not implement guarded edge writes".into(),
123        })
124    }
125    /// Observed form of [`GraphStore::upsert_edge_guarded`]. In addition to
126    /// the endpoint guard, it distinguishes create, live replacement, and
127    /// explicit resurrection without a caller-side read/write race.
128    async fn upsert_edge_guarded_observed(
129        &self,
130        _request: EdgeUpsertRequest,
131    ) -> StorageResult<GuardedEdgeUpsertOutcome> {
132        Err(StorageError::Unsupported {
133            capability: StorageCapability::Graph,
134            operation: "upsert_edge_guarded_observed".into(),
135            message: "this backend does not implement observed guarded edge upserts".into(),
136        })
137    }
138    /// Batch form of [`GraphStore::upsert_edge_guarded`]. All-or-nothing:
139    /// if any edge's endpoints are missing at write time, no edge from the
140    /// batch is persisted, `BatchWriteSummary::affected` is `0`, and
141    /// `GuardedBatchOutcome::refused` names the first failing batch entry and
142    /// its missing endpoint(s) — determined by the same in-transaction
143    /// pre-check that aborted the batch, not a post-hoc re-read.
144    /// Retain the original ordered input to enumerate every aborted write via
145    /// [`GuardedBatchOutcome::refusal_page`] beyond the default bounded sample.
146    ///
147    /// Default returns `StorageError::Unsupported`, for the same reason as
148    /// [`GraphStore::upsert_edge_guarded`]'s default.
149    async fn upsert_edges_guarded(&self, _edges: Vec<Edge>) -> StorageResult<GuardedBatchOutcome> {
150        Err(StorageError::Unsupported {
151            capability: StorageCapability::Graph,
152            operation: "upsert_edges_guarded".into(),
153            message: "this backend does not implement guarded edge writes".into(),
154        })
155    }
156    /// All-or-nothing observed batch form. Implementations must perform
157    /// endpoint and tombstone-policy preflight in the same write transaction
158    /// before applying any row.
159    async fn upsert_edges_guarded_observed(
160        &self,
161        _requests: Vec<EdgeUpsertRequest>,
162    ) -> StorageResult<GuardedEdgeBatchUpsertOutcome> {
163        Err(StorageError::Unsupported {
164            capability: StorageCapability::Graph,
165            operation: "upsert_edges_guarded_observed".into(),
166            message: "this backend does not implement observed guarded edge batches".into(),
167        })
168    }
169    /// Fetch an edge by link ID, returning `None` if absent. Filters soft-deleted rows.
170    async fn get_edge(&self, id: LinkId) -> StorageResult<Option<Edge>>;
171    /// Fetch an edge by link ID including soft-deleted rows. Used by the runtime hard-delete path
172    /// to locate and namespace-check an already-soft-deleted edge before purging it.
173    async fn get_edge_including_deleted(&self, id: LinkId) -> StorageResult<Option<Edge>>;
174    /// Fetch an edge by natural key (namespace, source, target, relation) including
175    /// soft-deleted rows. Used by the atomic-apply result renderer for a symmetric-relation
176    /// update whose surviving canonical row may be tombstoned (ADR-039 DO NOTHING) — the
177    /// normal `query_edges`/`list_edges` path filters `deleted_at IS NULL` and would report
178    /// "not found" for exactly that row.
179    ///
180    /// `namespace` is the natural key's own `namespace` column value (part of the
181    /// `UNIQUE(namespace, source_id, target_id, relation)` constraint this method queries by)
182    /// — it is passed explicitly rather than implied by whichever store instance `self` is,
183    /// so a caller who resolved the record's namespace independently of its own ambient token
184    /// (the atomic-apply renderer, which knows the committed edge's namespace from its prepare-
185    /// time `EdgeNaturalKey`, not from the caller's token) cannot accidentally query the wrong
186    /// namespace by relying on implicit store scoping.
187    async fn get_edge_by_natural_key_including_deleted(
188        &self,
189        namespace: &str,
190        source_id: Uuid,
191        target_id: Uuid,
192        relation: EdgeRelation,
193    ) -> StorageResult<Option<Edge>>;
194    /// Delete an edge by link ID using the specified delete mode.
195    async fn delete_edge(&self, id: LinkId, mode: DeleteMode) -> StorageResult<bool>;
196    /// Query edges with filter, sort, and pagination without an implicit
197    /// exact count. Implementations should return `total: None`; callers that
198    /// need a count use [`Self::count_edges`] explicitly.
199    async fn query_edges(
200        &self,
201        filter: EdgeFilter,
202        sort: Vec<SortOrder<EdgeSortField>>,
203        page: PageRequest,
204    ) -> StorageResult<Page<Edge>>;
205    /// Query edges across the given namespaces in one deterministic query
206    /// with real SQL paging. The multi-namespace analogue of
207    /// [`Self::query_edges`]: a single statement with `namespace IN (...)`
208    /// keeps `offset` continuation coherent, where fetching per-namespace
209    /// prefixes and slicing a client-side merge floats the window between
210    /// calls (silent duplicate/skip enumeration). Backends without batched
211    /// namespace support retain the single-namespace path and reject
212    /// multi-namespace requests explicitly. Implementations should return
213    /// `total: None`; callers that need a count use
214    /// [`Self::count_edges_in_namespaces`] explicitly.
215    async fn query_edges_in_namespaces(
216        &self,
217        namespaces: &[String],
218        filter: EdgeFilter,
219        sort: Vec<SortOrder<EdgeSortField>>,
220        page: PageRequest,
221    ) -> StorageResult<Page<Edge>> {
222        match namespaces.len() {
223            0 => Ok(Page {
224                items: Vec::new(),
225                total: None,
226            }),
227            1 => self.query_edges(filter, sort, page).await,
228            _ => Err(StorageError::Unsupported {
229                capability: StorageCapability::Graph,
230                operation: "query_edges_in_namespaces".into(),
231                message: "this backend does not implement batched namespace edge queries".into(),
232            }),
233        }
234    }
235    /// Count edges matching the given filter.
236    async fn count_edges(&self, filter: EdgeFilter) -> StorageResult<u64>;
237    /// Count edges across the given namespaces in one aggregate query.
238    /// Backends without batched namespace support retain the single-namespace
239    /// path and reject multi-namespace requests explicitly.
240    async fn count_edges_in_namespaces(
241        &self,
242        namespaces: &[String],
243        filter: EdgeFilter,
244    ) -> StorageResult<u64> {
245        match namespaces.len() {
246            0 => Ok(0),
247            1 => self.count_edges(filter).await,
248            _ => Err(StorageError::Unsupported {
249                capability: StorageCapability::Graph,
250                operation: "count_edges_in_namespaces".into(),
251                message: "this backend does not implement batched namespace edge counts".into(),
252            }),
253        }
254    }
255    /// Count edges grouped by relation, ignoring soft-deleted rows. Cheap
256    /// aggregate (`GROUP BY relation`) used to report the true per-relation
257    /// population for full-graph audits (#702.3).
258    async fn count_edges_by_relation(&self) -> StorageResult<Vec<(EdgeRelation, u64)>>;
259    /// Count edges grouped by relation across the given namespaces in one
260    /// aggregate query.
261    async fn count_edges_by_relation_in_namespaces(
262        &self,
263        namespaces: &[String],
264    ) -> StorageResult<Vec<(EdgeRelation, u64)>> {
265        match namespaces.len() {
266            0 => Ok(Vec::new()),
267            1 => self.count_edges_by_relation().await,
268            _ => Err(StorageError::Unsupported {
269                capability: StorageCapability::Graph,
270                operation: "count_edges_by_relation_in_namespaces".into(),
271                message: "this backend does not implement batched namespace relation counts".into(),
272            }),
273        }
274    }
275    /// Count live edges grouped by the base each endpoint resolves against.
276    ///
277    /// A relation breakdown cannot answer this: relations do not determine
278    /// endpoint bases, and the same relation appears on both sides of the
279    /// structure/provenance line. One aggregate query, same shape and cost as
280    /// the relation counts.
281    async fn count_edges_by_endpoint_base(&self) -> StorageResult<EdgeEndpointBaseCounts> {
282        Err(StorageError::Unsupported {
283            capability: StorageCapability::Graph,
284            operation: "count_edges_by_endpoint_base".into(),
285            message: "this backend does not implement endpoint-base edge counts".into(),
286        })
287    }
288    /// Count live edges by endpoint base across the given namespaces in one
289    /// aggregate query.
290    async fn count_edges_by_endpoint_base_in_namespaces(
291        &self,
292        namespaces: &[String],
293    ) -> StorageResult<EdgeEndpointBaseCounts> {
294        match namespaces.len() {
295            0 => Ok(EdgeEndpointBaseCounts::default()),
296            1 => self.count_edges_by_endpoint_base().await,
297            _ => Err(StorageError::Unsupported {
298                capability: StorageCapability::Graph,
299                operation: "count_edges_by_endpoint_base_in_namespaces".into(),
300                message: "this backend does not implement batched namespace endpoint-base counts"
301                    .into(),
302            }),
303        }
304    }
305    /// Seek-pagination page of edges ordered by `id` ascending, using an
306    /// indexed range scan (`id > after`) against the `(namespace, id)`
307    /// primary key instead of `OFFSET`. `after` is exclusive; `None` starts
308    /// from the beginning of the set. This remains an efficient compatibility
309    /// path for a fixed edge set, but random UUIDs inserted concurrently may
310    /// sort behind an issued boundary. Public concurrent walks use
311    /// [`Self::query_edges_sequence_after`] instead (#1424).
312    async fn query_edges_after(
313        &self,
314        filter: EdgeFilter,
315        after: Option<Uuid>,
316        limit: u32,
317    ) -> StorageResult<EdgeSeekPage>;
318    /// Resolve an edge id to its immutable insertion sequence.
319    async fn edge_sequence(&self, _id: Uuid) -> StorageResult<Option<i64>> {
320        Err(StorageError::Unsupported {
321            capability: StorageCapability::Graph,
322            operation: "edge_sequence".into(),
323            message: "this backend does not implement edge insertion sequences".into(),
324        })
325    }
326    /// Resolve edge ids to immutable insertion sequences. Implementations may
327    /// override this to batch the lookup; the default preserves correctness.
328    async fn edge_sequences(&self, ids: &[Uuid]) -> StorageResult<Vec<(Uuid, i64)>> {
329        let mut resolved = Vec::with_capacity(ids.len());
330        for id in ids {
331            if let Some(sequence) = self.edge_sequence(*id).await? {
332                resolved.push((*id, sequence));
333            }
334        }
335        Ok(resolved)
336    }
337    /// Seek-pagination page ordered by immutable insertion sequence. This is
338    /// the stable public-list contract for walks overlapping inserts (#1424).
339    async fn query_edges_sequence_after(
340        &self,
341        _filter: EdgeFilter,
342        _after: Option<SeekCursor>,
343        _limit: u32,
344    ) -> StorageResult<SeekPage<Edge>> {
345        Err(StorageError::Unsupported {
346            capability: StorageCapability::Graph,
347            operation: "query_edges_sequence_after".into(),
348            message: "this backend does not implement insertion-sequence edge pagination".into(),
349        })
350    }
351    /// Return immediate neighbors of a graph node.
352    async fn neighbors(
353        &self,
354        node_id: Uuid,
355        query: NeighborQuery,
356    ) -> StorageResult<Vec<NeighborHit>>;
357    /// Return one deterministic neighbor page. `after` is exclusive and
358    /// `neighbor_kinds`, when present, filters entity and note kinds before
359    /// the limit is applied. Backends that do not implement kind-aware paging
360    /// retain an explicit unsupported result rather than silently returning a
361    /// misleading page.
362    async fn neighbors_page(
363        &self,
364        node_id: Uuid,
365        mut query: NeighborQuery,
366        after: Option<NeighborCursor>,
367        neighbor_kinds: Option<Vec<String>>,
368    ) -> StorageResult<Vec<NeighborHit>> {
369        if neighbor_kinds
370            .as_ref()
371            .is_some_and(|kinds| !kinds.is_empty())
372        {
373            return Err(StorageError::Unsupported {
374                capability: StorageCapability::Graph,
375                operation: "neighbors_page".into(),
376                message: "this backend does not implement neighbor kind filtering".into(),
377            });
378        }
379        let limit = query.limit;
380        query.limit = None;
381        let mut hits = self.neighbors(node_id, query).await?;
382        if let Some(cursor) = after {
383            hits.retain(|hit| cursor.is_after(hit));
384        }
385        if let Some(limit) = limit {
386            hits.truncate(limit as usize);
387        }
388        Ok(hits)
389    }
390    /// Return neighbors in BOTH directions in a single call, each tagged with
391    /// the direction (`Out`/`In`) it was found in. `query.direction` is
392    /// ignored — this always fetches both directions.
393    ///
394    /// Exists so a caller that needs both-direction neighbors labeled by
395    /// direction (e.g. the `context` verb) can do so with one storage query
396    /// instead of two separate direction-scoped `neighbors` calls. The
397    /// default implementation preserves the original two-call behavior for
398    /// backends that don't override it; `SqlGraphStore` overrides this with a
399    /// single `UNION ALL` query that projects a direction literal per arm.
400    async fn neighbors_both_directions(
401        &self,
402        node_id: Uuid,
403        query: NeighborQuery,
404    ) -> StorageResult<Vec<DirectedNeighborHit>> {
405        let mut out_query = query.clone();
406        out_query.direction = Direction::Out;
407        let mut in_query = query;
408        in_query.direction = Direction::In;
409        let mut result = Vec::new();
410        for hit in self.neighbors(node_id, out_query).await? {
411            result.push(DirectedNeighborHit {
412                hit,
413                direction: Direction::Out,
414            });
415        }
416        for hit in self.neighbors(node_id, in_query).await? {
417            result.push(DirectedNeighborHit {
418                hit,
419                direction: Direction::In,
420            });
421        }
422        Ok(result)
423    }
424    /// Fetch multiple edges by their link IDs in a single round-trip.
425    ///
426    /// IDs that are not found (absent or soft-deleted) are silently skipped;
427    /// the returned `Vec` may be shorter than `ids`. Backends that support
428    /// batched `IN (...)` queries should override this; the default loops
429    /// `get_edge` so non-SQLite backends keep compiling unchanged.
430    ///
431    /// Callers must chunk large ID lists before calling if they need a strict
432    /// size bound; this method does not enforce a maximum.
433    async fn get_edges(&self, ids: &[LinkId]) -> StorageResult<Vec<Edge>> {
434        let mut out = Vec::with_capacity(ids.len());
435        for &id in ids {
436            if let Some(edge) = self.get_edge(id).await? {
437                out.push(edge);
438            }
439        }
440        Ok(out)
441    }
442    /// Return neighbors for multiple source nodes in a single round-trip,
443    /// yielding `(source_id, hit)` pairs.
444    ///
445    /// The `query` parameters (direction, relations, min_weight) are applied
446    /// uniformly to every source node. `query.limit` is applied **per source**:
447    /// each source returns at most `limit` hits. Backends that support batched
448    /// `source_id IN (...)` queries should override this; the default loops
449    /// `neighbors` so non-SQLite backends keep compiling unchanged.
450    async fn batch_neighbors(
451        &self,
452        sources: &[Uuid],
453        query: NeighborQuery,
454    ) -> StorageResult<Vec<(Uuid, NeighborHit)>> {
455        let mut out = Vec::new();
456        for &src in sources {
457            let hits = self.neighbors(src, query.clone()).await?;
458            for hit in hits {
459                out.push((src, hit));
460            }
461        }
462        Ok(out)
463    }
464    /// Bounded multi-hop BFS traversal from the given roots.
465    ///
466    /// Implementations must validate [`TraversalRequest::validate`], count
467    /// adjacency rows before first-visit de-duplication against the request's
468    /// shared execution budget, stop a root as soon as its effective result
469    /// limit is filled, and return an error rather than partial paths when the
470    /// work or time budget expires. Minimum-depth BFS selection is normative;
471    /// same-depth tie ordering is not.
472    async fn traverse(&self, request: TraversalRequest) -> StorageResult<Vec<GraphPath>>;
473    /// Hard-delete every incident edge (source or target) for `node_id`, regardless of soft-delete
474    /// state. Used during endpoint hard-delete to prevent dangling `graph_edges` rows (ADR-002
475    /// no-dangling-references contract).
476    async fn purge_incident_edges(&self, node_id: Uuid) -> StorageResult<u64>;
477}