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