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}