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}