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