Skip to main content

relay_knowledge/storage/sqlite/store/
implementations.rs

1//! Maps storage-port traits to the focused SQLite persistence owners.
2
3use std::{sync::Arc, time::Instant};
4
5use crate::{
6    domain::{
7        CodeChunkRecord, CodeGraphBatch, CodeGraphCommitReceipt, CodeReferenceRecord,
8        CodeSymbolRecord, CommitReceipt, GraphMutationBatch, GraphVersion, IndexKind, IndexStatus,
9    },
10    storage::{
11        AuditQueryRequest, CodeChunkSearchRequest, CodeGraphStore, CodeReferenceSearchRequest,
12        CodeSymbolSearchRequest, FileContentSearchHit, FileContentSearchRequest,
13        FileIndexDiagnostics, FileIndexRoot, FileIndexRootStatus, FileIndexRootUpdate,
14        FileSearchHit, FileSearchRequest, GraphCanvasStorageRequest, GraphCanvasStorageSnapshot,
15        GraphInspection, GraphSearchOutcome, GraphSearchRequest, GraphStore, HealthStorageSnapshot,
16        IndexCursor, IndexRefreshClaimRequest, IndexRefreshCompletion, IndexRefreshDiagnostics,
17        IndexRefreshFailure, IndexRefreshQueueRequest, IndexRefreshTask, IndexStore,
18        MutationLogEntry, MutationLogStore, NewAuditEvent, NewProposal, ProposalDecision,
19        ProposalListRequest, ServiceOperatorUpdate, StorageFuture, WorkerTaskClaimRequest,
20        WorkerTaskCompletion, WorkerTaskFailure, WorkerTaskSeed,
21    },
22};
23
24use super::super::{
25    canvas,
26    code::lifecycle::report,
27    code_graph, file_index,
28    graph::{commit_batch, current_graph_version, inspect_graph},
29    indexing,
30    mutation_log::read_mutations_after,
31    operations, retrieval,
32};
33use super::SqliteGraphStore;
34
35impl GraphStore for SqliteGraphStore {
36    fn commit_mutation_batch(&self, batch: GraphMutationBatch) -> StorageFuture<'_, CommitReceipt> {
37        self.run(move |connection| commit_batch(connection, batch))
38    }
39
40    fn inspect_graph(&self) -> StorageFuture<'_, GraphInspection> {
41        let database_path = self.database_path.clone();
42        let maintenance = Arc::clone(&self.maintenance);
43        self.run_read(move |connection| {
44            inspect_graph(connection, database_path.as_deref(), &maintenance)
45        })
46    }
47
48    fn health_snapshot(&self, now_ms: u64) -> StorageFuture<'_, HealthStorageSnapshot> {
49        let database_path = self.database_path.clone();
50        let maintenance = Arc::clone(&self.maintenance);
51        self.try_run_read(move |connection| {
52            Ok(HealthStorageSnapshot {
53                graph: inspect_graph(connection, database_path.as_deref(), &maintenance)?,
54                repository_code_totals: report::repository_totals(connection)?,
55                indexes: indexing::index_statuses(connection)?,
56                index_cursors: indexing::index_cursors(connection)?,
57                index_refresh: indexing::diagnostics(connection, now_ms)?,
58                file_index: file_index::diagnostics(connection)?,
59            })
60        })
61    }
62
63    fn graph_canvas(
64        &self,
65        request: GraphCanvasStorageRequest,
66    ) -> StorageFuture<'_, GraphCanvasStorageSnapshot> {
67        self.run_read(move |connection| canvas::graph_canvas(connection, request))
68    }
69
70    fn search(&self, request: GraphSearchRequest) -> StorageFuture<'_, GraphSearchOutcome> {
71        self.run_read(move |connection| retrieval::search_graph(connection, request))
72    }
73
74    fn current_graph_version(&self) -> StorageFuture<'_, GraphVersion> {
75        self.run_read(current_graph_version)
76    }
77}
78
79impl MutationLogStore for SqliteGraphStore {
80    fn read_after(
81        &self,
82        graph_version: GraphVersion,
83        limit: usize,
84    ) -> StorageFuture<'_, Vec<MutationLogEntry>> {
85        self.run_read(move |connection| read_mutations_after(connection, graph_version, limit))
86    }
87}
88
89impl IndexStore for SqliteGraphStore {
90    fn index_statuses(&self) -> StorageFuture<'_, Vec<IndexStatus>> {
91        self.run_read(|connection| indexing::index_statuses(connection))
92    }
93
94    fn mark_refresh_complete(
95        &self,
96        kind: IndexKind,
97        graph_version: GraphVersion,
98    ) -> StorageFuture<'_, IndexStatus> {
99        self.run(move |connection| indexing::mark_refresh_complete(connection, kind, graph_version))
100    }
101
102    fn index_cursors(&self) -> StorageFuture<'_, Vec<IndexCursor>> {
103        self.run_read(indexing::index_cursors)
104    }
105
106    fn queue_index_refreshes(
107        &self,
108        request: IndexRefreshQueueRequest,
109    ) -> StorageFuture<'_, IndexRefreshDiagnostics> {
110        self.run(move |connection| indexing::queue_index_refreshes(connection, request))
111    }
112
113    fn claim_index_refresh_task(
114        &self,
115        request: IndexRefreshClaimRequest,
116    ) -> StorageFuture<'_, Option<IndexRefreshTask>> {
117        self.run(move |connection| indexing::claim_index_refresh_task(connection, request))
118    }
119
120    fn complete_index_refresh_task(
121        &self,
122        request: IndexRefreshCompletion,
123    ) -> StorageFuture<'_, IndexRefreshTask> {
124        self.run(move |connection| indexing::complete_index_refresh_task(connection, request))
125    }
126
127    fn fail_index_refresh_task(
128        &self,
129        request: IndexRefreshFailure,
130    ) -> StorageFuture<'_, IndexRefreshTask> {
131        self.run(move |connection| indexing::fail_index_refresh_task(connection, request))
132    }
133
134    fn index_refresh_diagnostics(&self, now_ms: u64) -> StorageFuture<'_, IndexRefreshDiagnostics> {
135        self.run_read(move |connection| indexing::diagnostics(connection, now_ms))
136    }
137
138    fn queue_worker_tasks(
139        &self,
140        tasks: Vec<WorkerTaskSeed>,
141    ) -> StorageFuture<'_, Vec<crate::domain::WorkerTaskRecord>> {
142        self.run(move |connection| operations::queue_worker_tasks(connection, tasks))
143    }
144
145    fn worker_statuses(&self) -> StorageFuture<'_, Vec<crate::domain::WorkerStatus>> {
146        self.run_read(|connection| operations::worker_statuses(connection))
147    }
148
149    fn claim_worker_task(
150        &self,
151        request: WorkerTaskClaimRequest,
152    ) -> StorageFuture<'_, Option<crate::domain::WorkerTaskRecord>> {
153        self.run(move |connection| operations::claim_worker_task(connection, request))
154    }
155
156    fn complete_worker_task(
157        &self,
158        request: WorkerTaskCompletion,
159    ) -> StorageFuture<'_, crate::domain::WorkerTaskRecord> {
160        self.run(move |connection| operations::complete_worker_task(connection, request))
161    }
162
163    fn fail_worker_task(
164        &self,
165        request: WorkerTaskFailure,
166    ) -> StorageFuture<'_, crate::domain::WorkerTaskRecord> {
167        self.run(move |connection| operations::fail_worker_task(connection, request))
168    }
169
170    fn insert_proposal(
171        &self,
172        proposal: NewProposal,
173    ) -> StorageFuture<'_, crate::domain::ProposalRecord> {
174        self.run(move |connection| operations::insert_proposal(connection, proposal))
175    }
176
177    fn list_proposals(
178        &self,
179        request: ProposalListRequest,
180    ) -> StorageFuture<'_, Vec<crate::domain::ProposalRecord>> {
181        self.run_read(move |connection| operations::list_proposals(connection, request))
182    }
183
184    fn proposal_count(
185        &self,
186        state: Option<crate::domain::ProposalState>,
187    ) -> StorageFuture<'_, usize> {
188        self.run_read(move |connection| operations::proposal_count(connection, state))
189    }
190
191    fn proposal_by_id(
192        &self,
193        proposal_id: String,
194    ) -> StorageFuture<'_, Option<crate::domain::ProposalRecord>> {
195        self.run_read(move |connection| operations::proposal_by_id(connection, &proposal_id))
196    }
197
198    fn proposal_conflicts(
199        &self,
200        proposal_id: String,
201    ) -> StorageFuture<'_, Vec<crate::domain::ProposalConflictRecord>> {
202        self.run_read(move |connection| operations::proposal_conflicts(connection, &proposal_id))
203    }
204
205    fn decide_proposal(
206        &self,
207        request: ProposalDecision,
208    ) -> StorageFuture<'_, crate::domain::ProposalRecord> {
209        self.run(move |connection| operations::decide_proposal(connection, request))
210    }
211
212    fn insert_audit_event(
213        &self,
214        event: NewAuditEvent,
215    ) -> StorageFuture<'_, crate::domain::AuditEventRecord> {
216        self.run(move |connection| operations::insert_audit_event(connection, event))
217    }
218
219    fn query_audit_events(
220        &self,
221        request: AuditQueryRequest,
222    ) -> StorageFuture<'_, Vec<crate::domain::AuditEventRecord>> {
223        self.run_read(move |connection| operations::query_audit_events(connection, request))
224    }
225
226    fn audit_event_count(&self) -> StorageFuture<'_, usize> {
227        self.run_read(|connection| operations::audit_event_count(connection))
228    }
229
230    fn service_operator_status(&self) -> StorageFuture<'_, crate::domain::ServiceOperatorStatus> {
231        self.run_read(|connection| operations::service_operator_status(connection))
232    }
233
234    fn update_service_operator(
235        &self,
236        request: ServiceOperatorUpdate,
237    ) -> StorageFuture<'_, crate::domain::ServiceOperatorStatus> {
238        self.run(move |connection| operations::update_service_operator(connection, request))
239    }
240
241    fn replace_file_index_root(
242        &self,
243        update: FileIndexRootUpdate,
244    ) -> StorageFuture<'_, FileIndexRootStatus> {
245        self.run(move |connection| file_index::replace_root(connection, update))
246    }
247
248    fn mark_file_index_roots_unconfigured(
249        &self,
250        active_roots: Vec<FileIndexRoot>,
251        now_ms: u64,
252    ) -> StorageFuture<'_, FileIndexDiagnostics> {
253        self.run(move |connection| {
254            file_index::mark_unconfigured_roots(connection, active_roots, now_ms)
255        })
256    }
257
258    fn search_files(&self, request: FileSearchRequest) -> StorageFuture<'_, Vec<FileSearchHit>> {
259        let started = Instant::now();
260        let deadline = started
261            .checked_add(std::time::Duration::from_millis(request.timeout_ms))
262            .unwrap_or(started);
263        self.run_read_until(
264            deadline,
265            "file query timed out waiting for storage lock",
266            move |connection| file_index::search(connection, request, deadline),
267        )
268    }
269
270    fn search_file_content(
271        &self,
272        request: FileContentSearchRequest,
273    ) -> StorageFuture<'_, Vec<FileContentSearchHit>> {
274        let started = Instant::now();
275        let deadline = started
276            .checked_add(std::time::Duration::from_millis(request.timeout_ms))
277            .unwrap_or(started);
278        self.run_read_until(
279            deadline,
280            "file content query timed out waiting for storage lock",
281            move |connection| file_index::content::search(connection, request, deadline),
282        )
283    }
284
285    fn file_index_diagnostics(&self) -> StorageFuture<'_, FileIndexDiagnostics> {
286        self.run_read(|connection| file_index::diagnostics(connection))
287    }
288}
289
290impl CodeGraphStore for SqliteGraphStore {
291    fn commit_code_graph_batch(
292        &self,
293        batch: CodeGraphBatch,
294    ) -> StorageFuture<'_, CodeGraphCommitReceipt> {
295        self.run(move |connection| code_graph::commit_batch(connection, batch))
296    }
297
298    fn search_code_symbols(
299        &self,
300        request: CodeSymbolSearchRequest,
301    ) -> StorageFuture<'_, Vec<CodeSymbolRecord>> {
302        self.run_read(move |connection| code_graph::search_symbols(connection, request))
303    }
304
305    fn search_code_references(
306        &self,
307        request: CodeReferenceSearchRequest,
308    ) -> StorageFuture<'_, Vec<CodeReferenceRecord>> {
309        self.run_read(move |connection| code_graph::search_references(connection, request))
310    }
311
312    fn search_code_chunks(
313        &self,
314        request: CodeChunkSearchRequest,
315    ) -> StorageFuture<'_, Vec<CodeChunkRecord>> {
316        self.run_read(move |connection| code_graph::search_chunks(connection, request))
317    }
318}