Skip to main content

relay_knowledge/storage/sqlite/
store_impls.rs

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