Skip to main content

relay_knowledge/storage/sqlite/
store_impls.rs

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