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}