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