1use 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}