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