Skip to main content

khive_storage/
graph.rs

1//! Graph storage capability — edge CRUD and traversal.
2
3use async_trait::async_trait;
4use khive_types::EdgeRelation;
5use uuid::Uuid;
6
7use crate::capability::StorageCapability;
8use crate::error::StorageError;
9use crate::types::{
10    BatchWriteSummary, DeleteMode, DirectedNeighborHit, Direction, Edge, EdgeFilter, EdgeSeekPage,
11    EdgeSortField, GraphPath, GuardedBatchOutcome, GuardedWriteOutcome, LinkId, NeighborHit,
12    NeighborQuery, Page, PageRequest, SeekCursor, SeekPage, SortOrder, StorageResult,
13    TraversalRequest,
14};
15
16/// Directed edge CRUD and graph traversal over the knowledge graph.
17#[async_trait]
18pub trait GraphStore: Send + Sync + 'static {
19    /// Insert or update a single edge.
20    async fn upsert_edge(&self, edge: Edge) -> StorageResult<()>;
21    /// Insert or update a batch of edges.
22    async fn upsert_edges(&self, edges: Vec<Edge>) -> StorageResult<BatchWriteSummary>;
23    /// Insert or update a single edge, re-checking that both endpoints still
24    /// exist (and are not soft-deleted) as part of the same write, not a
25    /// separate prior read. Closes the TOCTOU window between an async
26    /// prepare-time existence check and a later, unconditional write: a
27    /// concurrent hard-delete of an endpoint that lands between the two can
28    /// otherwise leave a durably dangling edge (#769).
29    ///
30    /// Returns [`GuardedWriteOutcome::Refused`] naming exactly which
31    /// endpoint(s) were missing, determined by the guard's own in-transaction
32    /// probe — never reconstructed by a caller re-reading the endpoints after
33    /// the write already failed, since a concurrent write landing between the
34    /// refusal and any such later read could misreport which endpoint was
35    /// actually missing at write time.
36    ///
37    /// Default returns `StorageError::Unsupported`: a backend that does not
38    /// override this method cannot honor the endpoint-existence guarantee,
39    /// and silently falling back to [`GraphStore::upsert_edge`] would
40    /// reintroduce the TOCTOU window this method exists to close.
41    async fn upsert_edge_guarded(&self, _edge: Edge) -> StorageResult<GuardedWriteOutcome> {
42        Err(StorageError::Unsupported {
43            capability: StorageCapability::Graph,
44            operation: "upsert_edge_guarded".into(),
45            message: "this backend does not implement guarded edge writes".into(),
46        })
47    }
48    /// Batch form of [`GraphStore::upsert_edge_guarded`]. All-or-nothing:
49    /// if any edge's endpoints are missing at write time, no edge from the
50    /// batch is persisted, `BatchWriteSummary::affected` is `0`, and
51    /// `GuardedBatchOutcome::refused` names the first failing batch entry and
52    /// its missing endpoint(s) — determined by the same in-transaction
53    /// pre-check that aborted the batch, not a post-hoc re-read.
54    ///
55    /// Default returns `StorageError::Unsupported`, for the same reason as
56    /// [`GraphStore::upsert_edge_guarded`]'s default.
57    async fn upsert_edges_guarded(&self, _edges: Vec<Edge>) -> StorageResult<GuardedBatchOutcome> {
58        Err(StorageError::Unsupported {
59            capability: StorageCapability::Graph,
60            operation: "upsert_edges_guarded".into(),
61            message: "this backend does not implement guarded edge writes".into(),
62        })
63    }
64    /// Fetch an edge by link ID, returning `None` if absent. Filters soft-deleted rows.
65    async fn get_edge(&self, id: LinkId) -> StorageResult<Option<Edge>>;
66    /// Fetch an edge by link ID including soft-deleted rows. Used by the runtime hard-delete path
67    /// to locate and namespace-check an already-soft-deleted edge before purging it.
68    async fn get_edge_including_deleted(&self, id: LinkId) -> StorageResult<Option<Edge>>;
69    /// Fetch an edge by natural key (namespace, source, target, relation) including
70    /// soft-deleted rows. Used by the atomic-apply result renderer for a symmetric-relation
71    /// update whose surviving canonical row may be tombstoned (ADR-039 DO NOTHING) — the
72    /// normal `query_edges`/`list_edges` path filters `deleted_at IS NULL` and would report
73    /// "not found" for exactly that row.
74    ///
75    /// `namespace` is the natural key's own `namespace` column value (part of the
76    /// `UNIQUE(namespace, source_id, target_id, relation)` constraint this method queries by)
77    /// — it is passed explicitly rather than implied by whichever store instance `self` is,
78    /// so a caller who resolved the record's namespace independently of its own ambient token
79    /// (the atomic-apply renderer, which knows the committed edge's namespace from its prepare-
80    /// time `EdgeNaturalKey`, not from the caller's token) cannot accidentally query the wrong
81    /// namespace by relying on implicit store scoping.
82    async fn get_edge_by_natural_key_including_deleted(
83        &self,
84        namespace: &str,
85        source_id: Uuid,
86        target_id: Uuid,
87        relation: EdgeRelation,
88    ) -> StorageResult<Option<Edge>>;
89    /// Delete an edge by link ID using the specified delete mode.
90    async fn delete_edge(&self, id: LinkId, mode: DeleteMode) -> StorageResult<bool>;
91    /// Query edges with filter, sort, and pagination.
92    async fn query_edges(
93        &self,
94        filter: EdgeFilter,
95        sort: Vec<SortOrder<EdgeSortField>>,
96        page: PageRequest,
97    ) -> StorageResult<Page<Edge>>;
98    /// Count edges matching the given filter.
99    async fn count_edges(&self, filter: EdgeFilter) -> StorageResult<u64>;
100    /// Count edges across the given namespaces in one aggregate query.
101    /// Backends without batched namespace support retain the single-namespace
102    /// path and reject multi-namespace requests explicitly.
103    async fn count_edges_in_namespaces(
104        &self,
105        namespaces: &[String],
106        filter: EdgeFilter,
107    ) -> StorageResult<u64> {
108        match namespaces.len() {
109            0 => Ok(0),
110            1 => self.count_edges(filter).await,
111            _ => Err(StorageError::Unsupported {
112                capability: StorageCapability::Graph,
113                operation: "count_edges_in_namespaces".into(),
114                message: "this backend does not implement batched namespace edge counts".into(),
115            }),
116        }
117    }
118    /// Count edges grouped by relation, ignoring soft-deleted rows. Cheap
119    /// aggregate (`GROUP BY relation`) used to report the true per-relation
120    /// population for full-graph audits (#702.3).
121    async fn count_edges_by_relation(&self) -> StorageResult<Vec<(EdgeRelation, u64)>>;
122    /// Count edges grouped by relation across the given namespaces in one
123    /// aggregate query.
124    async fn count_edges_by_relation_in_namespaces(
125        &self,
126        namespaces: &[String],
127    ) -> StorageResult<Vec<(EdgeRelation, u64)>> {
128        match namespaces.len() {
129            0 => Ok(Vec::new()),
130            1 => self.count_edges_by_relation().await,
131            _ => Err(StorageError::Unsupported {
132                capability: StorageCapability::Graph,
133                operation: "count_edges_by_relation_in_namespaces".into(),
134                message: "this backend does not implement batched namespace relation counts".into(),
135            }),
136        }
137    }
138    /// Seek-pagination page of edges ordered by `id` ascending, using an
139    /// indexed range scan (`id > after`) against the `(namespace, id)`
140    /// primary key instead of `OFFSET`. `after` is exclusive; `None` starts
141    /// from the beginning of the set. This remains an efficient compatibility
142    /// path for a fixed edge set, but random UUIDs inserted concurrently may
143    /// sort behind an issued boundary. Public concurrent walks use
144    /// [`Self::query_edges_sequence_after`] instead (#1424).
145    async fn query_edges_after(
146        &self,
147        filter: EdgeFilter,
148        after: Option<Uuid>,
149        limit: u32,
150    ) -> StorageResult<EdgeSeekPage>;
151    /// Resolve an edge id to its immutable insertion sequence.
152    async fn edge_sequence(&self, _id: Uuid) -> StorageResult<Option<i64>> {
153        Err(StorageError::Unsupported {
154            capability: StorageCapability::Graph,
155            operation: "edge_sequence".into(),
156            message: "this backend does not implement edge insertion sequences".into(),
157        })
158    }
159    /// Resolve edge ids to immutable insertion sequences. Implementations may
160    /// override this to batch the lookup; the default preserves correctness.
161    async fn edge_sequences(&self, ids: &[Uuid]) -> StorageResult<Vec<(Uuid, i64)>> {
162        let mut resolved = Vec::with_capacity(ids.len());
163        for id in ids {
164            if let Some(sequence) = self.edge_sequence(*id).await? {
165                resolved.push((*id, sequence));
166            }
167        }
168        Ok(resolved)
169    }
170    /// Seek-pagination page ordered by immutable insertion sequence. This is
171    /// the stable public-list contract for walks overlapping inserts (#1424).
172    async fn query_edges_sequence_after(
173        &self,
174        _filter: EdgeFilter,
175        _after: Option<SeekCursor>,
176        _limit: u32,
177    ) -> StorageResult<SeekPage<Edge>> {
178        Err(StorageError::Unsupported {
179            capability: StorageCapability::Graph,
180            operation: "query_edges_sequence_after".into(),
181            message: "this backend does not implement insertion-sequence edge pagination".into(),
182        })
183    }
184    /// Return immediate neighbors of a graph node.
185    async fn neighbors(
186        &self,
187        node_id: Uuid,
188        query: NeighborQuery,
189    ) -> StorageResult<Vec<NeighborHit>>;
190    /// Return neighbors in BOTH directions in a single call, each tagged with
191    /// the direction (`Out`/`In`) it was found in. `query.direction` is
192    /// ignored — this always fetches both directions.
193    ///
194    /// Exists so a caller that needs both-direction neighbors labeled by
195    /// direction (e.g. the `context` verb) can do so with one storage query
196    /// instead of two separate direction-scoped `neighbors` calls. The
197    /// default implementation preserves the original two-call behavior for
198    /// backends that don't override it; `SqlGraphStore` overrides this with a
199    /// single `UNION ALL` query that projects a direction literal per arm.
200    async fn neighbors_both_directions(
201        &self,
202        node_id: Uuid,
203        query: NeighborQuery,
204    ) -> StorageResult<Vec<DirectedNeighborHit>> {
205        let mut out_query = query.clone();
206        out_query.direction = Direction::Out;
207        let mut in_query = query;
208        in_query.direction = Direction::In;
209        let mut result = Vec::new();
210        for hit in self.neighbors(node_id, out_query).await? {
211            result.push(DirectedNeighborHit {
212                hit,
213                direction: Direction::Out,
214            });
215        }
216        for hit in self.neighbors(node_id, in_query).await? {
217            result.push(DirectedNeighborHit {
218                hit,
219                direction: Direction::In,
220            });
221        }
222        Ok(result)
223    }
224    /// Fetch multiple edges by their link IDs in a single round-trip.
225    ///
226    /// IDs that are not found (absent or soft-deleted) are silently skipped;
227    /// the returned `Vec` may be shorter than `ids`. Backends that support
228    /// batched `IN (...)` queries should override this; the default loops
229    /// `get_edge` so non-SQLite backends keep compiling unchanged.
230    ///
231    /// Callers must chunk large ID lists before calling if they need a strict
232    /// size bound; this method does not enforce a maximum.
233    async fn get_edges(&self, ids: &[LinkId]) -> StorageResult<Vec<Edge>> {
234        let mut out = Vec::with_capacity(ids.len());
235        for &id in ids {
236            if let Some(edge) = self.get_edge(id).await? {
237                out.push(edge);
238            }
239        }
240        Ok(out)
241    }
242    /// Return neighbors for multiple source nodes in a single round-trip,
243    /// yielding `(source_id, hit)` pairs.
244    ///
245    /// The `query` parameters (direction, relations, min_weight) are applied
246    /// uniformly to every source node. `query.limit` is applied **per source**:
247    /// each source returns at most `limit` hits. Backends that support batched
248    /// `source_id IN (...)` queries should override this; the default loops
249    /// `neighbors` so non-SQLite backends keep compiling unchanged.
250    async fn batch_neighbors(
251        &self,
252        sources: &[Uuid],
253        query: NeighborQuery,
254    ) -> StorageResult<Vec<(Uuid, NeighborHit)>> {
255        let mut out = Vec::new();
256        for &src in sources {
257            let hits = self.neighbors(src, query.clone()).await?;
258            for hit in hits {
259                out.push((src, hit));
260            }
261        }
262        Ok(out)
263    }
264    /// Bounded multi-hop BFS traversal from the given roots.
265    ///
266    /// Implementations must validate [`TraversalRequest::validate`], count
267    /// adjacency rows before first-visit de-duplication against the request's
268    /// shared execution budget, stop a root as soon as its effective result
269    /// limit is filled, and return an error rather than partial paths when the
270    /// work or time budget expires. Minimum-depth BFS selection is normative;
271    /// same-depth tie ordering is not.
272    async fn traverse(&self, request: TraversalRequest) -> StorageResult<Vec<GraphPath>>;
273    /// Hard-delete every incident edge (source or target) for `node_id`, regardless of soft-delete
274    /// state. Used during endpoint hard-delete to prevent dangling `graph_edges` rows (ADR-002
275    /// no-dangling-references contract).
276    async fn purge_incident_edges(&self, node_id: Uuid) -> StorageResult<u64>;
277}