Skip to main content

relay_knowledge/storage/sqlite/
store_impls.rs

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