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