1use std::{path::Path, sync::Arc};
2
3#[path = "partitioned/catalog.rs"]
4mod catalog;
5#[path = "partitioned/control_delegates.rs"]
6mod control_delegates;
7#[path = "partitioned/diagnostics.rs"]
8mod diagnostics;
9#[path = "partitioned/retention.rs"]
10mod retention;
11#[path = "partitioned/routing.rs"]
12mod routing;
13#[path = "partitioned/status.rs"]
14mod status;
15#[path = "partitioned/totals.rs"]
16mod totals;
17
18use crate::{
19 domain::{
20 CodeFeatureFlagGraph, CodeFeatureFlagRequest, CodeIndexBatch, CodeIndexCheckpoint,
21 CodeIndexSession, CodeIndexSnapshot, CodeIndexSummary, CodeRepositoryCrossEdge,
22 CodeRepositoryRegistration, CodeRepositoryRemovalSummary, CodeRepositoryReport,
23 CodeRepositorySet, CodeRepositorySetMember, CodeRepositorySetRefreshSummary,
24 CodeRepositorySetStatus, CodeRepositoryStatus, CodeRepositoryTotals, CodeRetrievalHit,
25 CodeRetrievalRequest, CodeSymbolGenerationCounts, SoftwareGlobalProjection,
26 SoftwareGlobalRequest,
27 },
28 paths::RuntimePaths,
29 storage::{
30 CodeImpactChanges, CodeIndexTaskClaimRequest, CodeIndexTaskCompletion,
31 CodeIndexTaskFailure, CodeIndexTaskLeaseRecord, CodeIndexTaskLeaseRecovery,
32 CodeIndexTaskLeaseRenewal, CodeRepositorySetMemberSeed,
33 CodeRepositorySetRefreshTaskClaimRequest, CodeRepositorySetRefreshTaskCompletion,
34 CodeRepositorySetRefreshTaskFailure, CodeRepositorySetRefreshTaskSeed,
35 CodeRepositorySetSeed, CodeRepositoryStore, CodeScopeRetentionRequest, SqliteGraphStore,
36 StorageError, StorageFuture,
37 },
38};
39
40use catalog::{SqliteShardCatalog, initialize_catalog_schema};
41use retention::merge_scope_retention_summaries;
42use routing::{
43 current_control_scope, is_missing_code_scope_error, repository_store_for_selector,
44 source_scope_store,
45};
46use status::mirror_status;
47
48#[derive(Clone)]
51pub struct PartitionedSqliteKnowledgeStore {
52 control: Arc<SqliteGraphStore>,
53 catalog: Arc<SqliteShardCatalog>,
54}
55
56impl PartitionedSqliteKnowledgeStore {
57 pub fn open(control_path: impl AsRef<Path>, paths: RuntimePaths) -> Result<Self, StorageError> {
58 let control_path = control_path.as_ref().to_path_buf();
59 let control = Arc::new(SqliteGraphStore::open(&control_path)?);
60 initialize_catalog_schema(&control_path)?;
61
62 Ok(Self {
63 control,
64 catalog: Arc::new(SqliteShardCatalog::new(control_path, paths)),
65 })
66 }
67}
68
69impl CodeRepositoryStore for PartitionedSqliteKnowledgeStore {
70 fn upsert_code_repository(
71 &self,
72 registration: CodeRepositoryRegistration,
73 ) -> StorageFuture<'_, CodeRepositoryStatus> {
74 let this = self.clone();
75 Box::pin(async move {
76 let status = this
77 .control
78 .upsert_code_repository(registration.clone())
79 .await?;
80 let imported_scope = status.last_indexed_scope_id.clone();
81 let shard = this
82 .catalog
83 .staged_repository_store(status.repository_id.clone())
84 .await?;
85 this.catalog
86 .import_control_repository(
87 Arc::clone(&shard),
88 status.repository_id.clone(),
89 imported_scope.clone(),
90 )
91 .await?;
92 let shard_status = shard.upsert_code_repository(registration).await?;
93 if let Some(source_scope) = imported_scope {
94 this.catalog
95 .record_scope(status.repository_id.clone(), source_scope)
96 .await?;
97 } else {
98 this.catalog
99 .activate_repository(status.repository_id.clone())
100 .await?;
101 }
102 Ok(CodeRepositoryStatus {
103 alias: status.alias,
104 ..shard_status
105 })
106 })
107 }
108
109 fn code_repository_status(
110 &self,
111 repository: String,
112 ) -> StorageFuture<'_, Option<CodeRepositoryStatus>> {
113 let this = self.clone();
114 Box::pin(async move {
115 let Some(control_status) = this.control.code_repository_status(repository).await?
116 else {
117 return Ok(None);
118 };
119 let Some(shard) = this
120 .catalog
121 .existing_repository_store(control_status.repository_id.clone())
122 .await?
123 else {
124 return Ok(Some(control_status));
125 };
126 let Some(mut shard_status) = shard
127 .code_repository_status(control_status.repository_id.clone())
128 .await?
129 else {
130 return Ok(Some(control_status));
131 };
132 shard_status.alias = control_status.alias;
133 Ok(Some(shard_status))
134 })
135 }
136
137 fn list_code_repositories(&self) -> StorageFuture<'_, Vec<CodeRepositoryStatus>> {
138 control_delegates::list_code_repositories(self)
139 }
140
141 fn remove_code_repository(
142 &self,
143 repository: String,
144 now_ms: u64,
145 ) -> StorageFuture<'_, Option<CodeRepositoryRemovalSummary>> {
146 let this = self.clone();
147 Box::pin(async move {
148 let Some(control_status) = this.control.code_repository_status(repository).await?
149 else {
150 return Ok(None);
151 };
152 let shard = this
153 .catalog
154 .existing_repository_store(control_status.repository_id.clone())
155 .await?;
156 let removed = this
157 .control
158 .remove_code_repository(control_status.repository_id.clone(), now_ms)
159 .await?;
160 let Some(summary) = removed else {
161 return Ok(None);
162 };
163 if let Some(shard) = shard {
164 shard
165 .remove_code_repository(control_status.repository_id.clone(), now_ms)
166 .await?;
167 }
168 this.catalog
169 .remove_repository(control_status.repository_id)
170 .await?;
171 Ok(Some(summary))
172 })
173 }
174
175 fn code_repository_scope_status(
176 &self,
177 repository: String,
178 resolved_commit_sha: String,
179 path_filters: Vec<String>,
180 language_filters: Vec<String>,
181 ) -> StorageFuture<'_, Option<CodeRepositoryStatus>> {
182 let this = self.clone();
183 Box::pin(async move {
184 let Some(control_status) = this.control.code_repository_status(repository).await?
185 else {
186 return Ok(None);
187 };
188 let Some(shard) = this
189 .catalog
190 .existing_repository_store(control_status.repository_id.clone())
191 .await?
192 else {
193 return this
194 .control
195 .code_repository_scope_status(
196 control_status.repository_id,
197 resolved_commit_sha,
198 path_filters,
199 language_filters,
200 )
201 .await;
202 };
203 let status = shard
204 .code_repository_scope_status(
205 control_status.repository_id.clone(),
206 resolved_commit_sha.clone(),
207 path_filters.clone(),
208 language_filters.clone(),
209 )
210 .await?;
211 if let Some(mut status) = status {
212 status.alias = control_status.alias;
213 return Ok(Some(status));
214 }
215 this.control
216 .code_repository_scope_status(
217 control_status.repository_id,
218 resolved_commit_sha,
219 path_filters,
220 language_filters,
221 )
222 .await
223 })
224 }
225
226 fn latest_code_repository_scope_status(
227 &self,
228 repository: String,
229 path_filters: Vec<String>,
230 language_filters: Vec<String>,
231 ) -> StorageFuture<'_, Option<CodeRepositoryStatus>> {
232 let this = self.clone();
233 Box::pin(async move {
234 let Some(control_status) = this.control.code_repository_status(repository).await?
235 else {
236 return Ok(None);
237 };
238 let Some(shard) = this
239 .catalog
240 .existing_repository_store(control_status.repository_id.clone())
241 .await?
242 else {
243 return this
244 .control
245 .latest_code_repository_scope_status(
246 control_status.repository_id,
247 path_filters,
248 language_filters,
249 )
250 .await;
251 };
252 let status = shard
253 .latest_code_repository_scope_status(
254 control_status.repository_id.clone(),
255 path_filters.clone(),
256 language_filters.clone(),
257 )
258 .await?;
259 if let Some(mut status) = status {
260 status.alias = control_status.alias;
261 return Ok(Some(status));
262 }
263 this.control
264 .latest_code_repository_scope_status(
265 control_status.repository_id,
266 path_filters,
267 language_filters,
268 )
269 .await
270 })
271 }
272
273 fn queue_code_index_task(
274 &self,
275 task: crate::storage::CodeIndexTaskSeed,
276 ) -> StorageFuture<'_, crate::domain::CodeIndexTaskRecord> {
277 let control = Arc::clone(&self.control);
278 Box::pin(async move { control.queue_code_index_task(task).await })
279 }
280
281 fn claim_code_index_task(
282 &self,
283 request: CodeIndexTaskClaimRequest,
284 ) -> StorageFuture<'_, Option<crate::domain::CodeIndexTaskRecord>> {
285 self.control.claim_code_index_task(request)
286 }
287
288 fn recover_code_index_task_leases(
289 &self,
290 now_ms: u64,
291 max_attempts: u32,
292 ) -> StorageFuture<'_, ()> {
293 self.control
294 .recover_code_index_task_leases(now_ms, max_attempts)
295 }
296
297 fn running_code_index_task_leases(&self) -> StorageFuture<'_, Vec<CodeIndexTaskLeaseRecord>> {
298 self.control.running_code_index_task_leases()
299 }
300
301 fn recover_code_index_task_leases_by_task(
302 &self,
303 request: CodeIndexTaskLeaseRecovery,
304 ) -> StorageFuture<'_, usize> {
305 self.control.recover_code_index_task_leases_by_task(request)
306 }
307
308 fn reset_code_index_tasks(
309 &self,
310 repository_id: String,
311 now_ms: u64,
312 ) -> StorageFuture<'_, Vec<crate::domain::CodeIndexTaskRecord>> {
313 self.control.reset_code_index_tasks(repository_id, now_ms)
314 }
315
316 fn renew_code_index_task_lease(
317 &self,
318 request: CodeIndexTaskLeaseRenewal,
319 ) -> StorageFuture<'_, crate::domain::CodeIndexTaskRecord> {
320 self.control.renew_code_index_task_lease(request)
321 }
322
323 fn complete_code_index_task(
324 &self,
325 request: CodeIndexTaskCompletion,
326 ) -> StorageFuture<'_, crate::domain::CodeIndexTaskRecord> {
327 self.control.complete_code_index_task(request)
328 }
329
330 fn fail_code_index_task(
331 &self,
332 request: CodeIndexTaskFailure,
333 ) -> StorageFuture<'_, crate::domain::CodeIndexTaskRecord> {
334 self.control.fail_code_index_task(request)
335 }
336
337 fn code_index_task(
338 &self,
339 task_id: String,
340 ) -> StorageFuture<'_, Option<crate::domain::CodeIndexTaskRecord>> {
341 self.control.code_index_task(task_id)
342 }
343
344 fn active_code_index_task(
345 &self,
346 repository_id: String,
347 ) -> StorageFuture<'_, Option<crate::domain::CodeIndexTaskRecord>> {
348 self.control.active_code_index_task(repository_id)
349 }
350
351 fn code_index_task_queue_status(
352 &self,
353 ) -> StorageFuture<'_, crate::domain::CodeIndexTaskQueueStatus> {
354 self.control.code_index_task_queue_status()
355 }
356 fn code_index_checkpoint(
357 &self,
358 source_scope: String,
359 ) -> StorageFuture<'_, Option<CodeIndexCheckpoint>> {
360 let this = self.clone();
361 Box::pin(async move {
362 if let Some(shard) = this
363 .catalog
364 .checkpoint_scope_store(source_scope.clone())
365 .await?
366 {
367 if let Some(checkpoint) = shard.code_index_checkpoint(source_scope.clone()).await? {
368 return Ok(Some(checkpoint));
369 }
370 }
371 this.control.code_index_checkpoint(source_scope).await
372 })
373 }
374
375 fn latest_code_index_checkpoint(
376 &self,
377 repository_id: String,
378 ) -> StorageFuture<'_, Option<CodeIndexCheckpoint>> {
379 let this = self.clone();
380 Box::pin(async move {
381 if let Some(shard) = this
382 .catalog
383 .checkpoint_repository_store(repository_id.clone())
384 .await?
385 {
386 return shard.latest_code_index_checkpoint(repository_id).await;
387 }
388 this.control
389 .latest_code_index_checkpoint(repository_id)
390 .await
391 })
392 }
393
394 fn code_scope_retention(
395 &self,
396 repository_id: String,
397 ) -> StorageFuture<'_, crate::domain::CodeScopeRetentionSummary> {
398 let this = self.clone();
399 Box::pin(async move {
400 if let Some(shard) = this
401 .catalog
402 .existing_repository_store(repository_id.clone())
403 .await?
404 {
405 return shard.code_scope_retention(repository_id).await;
406 }
407 this.control.code_scope_retention(repository_id).await
408 })
409 }
410
411 fn prune_code_repository_scopes(
412 &self,
413 request: CodeScopeRetentionRequest,
414 ) -> StorageFuture<'_, crate::domain::CodeScopeRetentionSummary> {
415 let this = self.clone();
416 Box::pin(async move {
417 if let Some(shard) = this
418 .catalog
419 .existing_repository_store(request.repository_id.clone())
420 .await?
421 {
422 let control_retention = this
423 .control
424 .prune_code_repository_scopes(request.clone())
425 .await?;
426 let shard_retention = shard
427 .prune_code_repository_scopes_with_retained(
428 request.clone(),
429 control_retention.retained_scopes.clone(),
430 )
431 .await;
432 return shard_retention.map(|summary| {
433 merge_scope_retention_summaries(
434 request.repository_id,
435 control_retention,
436 summary,
437 )
438 });
439 }
440 this.control.prune_code_repository_scopes(request).await
441 })
442 }
443
444 fn code_file_fingerprints(
445 &self,
446 repository_id: String,
447 ) -> StorageFuture<'_, Vec<crate::domain::CodeFileFingerprint>> {
448 let this = self.clone();
449 Box::pin(async move {
450 if let Some(shard) = this
451 .catalog
452 .existing_repository_store(repository_id.clone())
453 .await?
454 {
455 return shard.code_file_fingerprints(repository_id).await;
456 }
457 this.control.code_file_fingerprints(repository_id).await
458 })
459 }
460
461 fn code_file_fingerprints_for_scope(
462 &self,
463 source_scope: String,
464 ) -> StorageFuture<'_, Vec<crate::domain::CodeFileFingerprint>> {
465 let this = self.clone();
466 Box::pin(async move {
467 if let Some(shard) = source_scope_store(&this.catalog, source_scope.clone()).await? {
468 return shard.code_file_fingerprints_for_scope(source_scope).await;
469 }
470 this.control
471 .code_file_fingerprints_for_scope(source_scope)
472 .await
473 })
474 }
475
476 fn code_file_candidate_paths_for_scope(
477 &self,
478 source_scope: String,
479 path_filters: Vec<String>,
480 language_filters: Vec<String>,
481 exclude_generated: bool,
482 limit: usize,
483 ) -> StorageFuture<'_, Vec<String>> {
484 let this = self.clone();
485 Box::pin(async move {
486 if let Some(shard) = source_scope_store(&this.catalog, source_scope.clone()).await? {
487 return shard
488 .code_file_candidate_paths_for_scope(
489 source_scope,
490 path_filters,
491 language_filters,
492 exclude_generated,
493 limit,
494 )
495 .await;
496 }
497 this.control
498 .code_file_candidate_paths_for_scope(
499 source_scope,
500 path_filters,
501 language_filters,
502 exclude_generated,
503 limit,
504 )
505 .await
506 })
507 }
508
509 fn code_file_candidate_paths_for_query_scope(
510 &self,
511 source_scope: String,
512 query: String,
513 path_filters: Vec<String>,
514 language_filters: Vec<String>,
515 exclude_generated: bool,
516 limit: usize,
517 ) -> StorageFuture<'_, Vec<String>> {
518 let this = self.clone();
519 Box::pin(async move {
520 if let Some(shard) = source_scope_store(&this.catalog, source_scope.clone()).await? {
521 return shard
522 .code_file_candidate_paths_for_query_scope(
523 source_scope,
524 query,
525 path_filters,
526 language_filters,
527 exclude_generated,
528 limit,
529 )
530 .await;
531 }
532 this.control
533 .code_file_candidate_paths_for_query_scope(
534 source_scope,
535 query,
536 path_filters,
537 language_filters,
538 exclude_generated,
539 limit,
540 )
541 .await
542 })
543 }
544
545 fn apply_code_index_snapshot(
546 &self,
547 snapshot: CodeIndexSnapshot,
548 ) -> StorageFuture<'_, CodeIndexSummary> {
549 let this = self.clone();
550 Box::pin(async move {
551 let base_scope = control_delegates::incremental_base_scope(&this, &snapshot).await?;
552 let shard = if snapshot.full_replace {
553 this.catalog
554 .staged_repository_store(snapshot.repository_id.clone())
555 .await?
556 } else {
557 match this
558 .catalog
559 .existing_repository_store(snapshot.repository_id.clone())
560 .await?
561 {
562 Some(shard) => shard,
563 None => {
564 this.catalog
565 .staged_repository_store(snapshot.repository_id.clone())
566 .await?
567 }
568 }
569 };
570 this.catalog
571 .import_control_repository(
572 Arc::clone(&shard),
573 snapshot.repository_id.clone(),
574 base_scope,
575 )
576 .await?;
577 let summary = shard.apply_code_index_snapshot(snapshot).await?;
578 let status = shard
579 .code_repository_status(summary.repository_id.clone())
580 .await?
581 .ok_or_else(|| {
582 StorageError::InvalidInput(
583 "sharded code repository status is missing after index".to_owned(),
584 )
585 })?;
586 this.catalog
587 .record_scope(summary.repository_id.clone(), summary.source_scope.clone())
588 .await?;
589 mirror_status(&this.control, status).await?;
590 Ok(summary)
591 })
592 }
593
594 fn clear_code_workspace_state(
595 &self,
596 repository_id: String,
597 source_scope: String,
598 ) -> StorageFuture<'_, ()> {
599 let this = self.clone();
600 Box::pin(async move {
601 if let Some(shard) = this
602 .catalog
603 .existing_repository_store(repository_id.clone())
604 .await?
605 {
606 shard
607 .clear_code_workspace_state(repository_id.clone(), source_scope.clone())
608 .await?;
609 }
610 this.control
611 .clear_code_workspace_state(repository_id, source_scope)
612 .await
613 })
614 }
615 fn begin_code_index_session(
616 &self,
617 session: CodeIndexSession,
618 ) -> StorageFuture<'_, CodeIndexCheckpoint> {
619 let this = self.clone();
620 Box::pin(async move {
621 let repository_id = session.repository_id.clone();
622 let source_scope = session.source_scope.clone();
623 let shard = this
624 .catalog
625 .staged_repository_store(repository_id.clone())
626 .await?;
627 let control_scope = current_control_scope(&this.control, repository_id.clone()).await?;
628 this.catalog
629 .import_control_repository(Arc::clone(&shard), repository_id.clone(), control_scope)
630 .await?;
631 let checkpoint = shard.begin_code_index_session(session).await?;
632 this.catalog
633 .stage_scope(repository_id, source_scope)
634 .await?;
635 Ok(checkpoint)
636 })
637 }
638
639 fn apply_code_index_batch(
640 &self,
641 batch: CodeIndexBatch,
642 ) -> StorageFuture<'_, CodeIndexCheckpoint> {
643 let this = self.clone();
644 Box::pin(async move {
645 let repository_id = batch.repository_id.clone();
646 let source_scope = batch.source_scope.clone();
647 let shard = this
648 .catalog
649 .staged_repository_store(repository_id.clone())
650 .await?;
651 let control_scope = current_control_scope(&this.control, repository_id.clone()).await?;
652 this.catalog
653 .import_control_repository(Arc::clone(&shard), repository_id.clone(), control_scope)
654 .await?;
655 let checkpoint = shard.apply_code_index_batch(batch).await?;
656 this.catalog
657 .stage_scope(repository_id, source_scope)
658 .await?;
659 Ok(checkpoint)
660 })
661 }
662
663 fn finalize_code_index_session(
664 &self,
665 session: CodeIndexSession,
666 ) -> StorageFuture<'_, CodeIndexSummary> {
667 let this = self.clone();
668 Box::pin(async move {
669 let shard = this
670 .catalog
671 .staged_repository_store(session.repository_id.clone())
672 .await?;
673 let summary = shard.finalize_code_index_session(session).await?;
674 let status = shard
675 .code_repository_status(summary.repository_id.clone())
676 .await?
677 .ok_or_else(|| {
678 StorageError::InvalidInput(
679 "sharded code repository status is missing after finalize".to_owned(),
680 )
681 })?;
682 this.catalog
683 .record_scope(summary.repository_id.clone(), summary.source_scope.clone())
684 .await?;
685 mirror_status(&this.control, status).await?;
686 Ok(summary)
687 })
688 }
689
690 fn search_code(
691 &self,
692 request: CodeRetrievalRequest,
693 ) -> StorageFuture<'_, Vec<CodeRetrievalHit>> {
694 let this = self.clone();
695 Box::pin(async move {
696 if let Some(shard) = repository_store_for_selector(
697 &this.control,
698 &this.catalog,
699 request.repository.repository.clone(),
700 )
701 .await?
702 {
703 return match shard.search_code(request.clone()).await {
704 Ok(hits) => Ok(hits),
705 Err(error) if is_missing_code_scope_error(&error) => {
706 this.control.search_code(request).await
707 }
708 Err(error) => Err(error),
709 };
710 }
711 this.control.search_code(request).await
712 })
713 }
714
715 fn search_code_feature_flags(
716 &self,
717 request: CodeFeatureFlagRequest,
718 ) -> StorageFuture<'_, Vec<CodeFeatureFlagGraph>> {
719 let this = self.clone();
720 Box::pin(async move {
721 if let Some(shard) = repository_store_for_selector(
722 &this.control,
723 &this.catalog,
724 request.repository.repository.clone(),
725 )
726 .await?
727 {
728 return match shard.search_code_feature_flags(request.clone()).await {
729 Ok(flags) => Ok(flags),
730 Err(error) if is_missing_code_scope_error(&error) => {
731 this.control.search_code_feature_flags(request).await
732 }
733 Err(error) => Err(error),
734 };
735 }
736 this.control.search_code_feature_flags(request).await
737 })
738 }
739
740 fn search_code_feature_flags_scope(
741 &self,
742 source_scope: String,
743 request: CodeFeatureFlagRequest,
744 ) -> StorageFuture<'_, Vec<CodeFeatureFlagGraph>> {
745 routing::search_code_feature_flags_scope(self.clone(), source_scope, request)
746 }
747
748 fn search_code_scope(
749 &self,
750 source_scope: String,
751 request: CodeRetrievalRequest,
752 ) -> StorageFuture<'_, Vec<CodeRetrievalHit>> {
753 routing::search_code_scope(self.clone(), source_scope, request)
754 }
755
756 fn analyze_code_impact(
757 &self,
758 request: crate::domain::CodeImpactRequest,
759 changes: CodeImpactChanges,
760 ) -> StorageFuture<'_, Vec<CodeRetrievalHit>> {
761 let this = self.clone();
762 Box::pin(async move {
763 if let Some(shard) = repository_store_for_selector(
764 &this.control,
765 &this.catalog,
766 request.repository.repository.clone(),
767 )
768 .await?
769 {
770 return match shard
771 .analyze_code_impact(request.clone(), changes.clone())
772 .await
773 {
774 Ok(hits) => Ok(hits),
775 Err(error) if is_missing_code_scope_error(&error) => {
776 this.control.analyze_code_impact(request, changes).await
777 }
778 Err(error) => Err(error),
779 };
780 }
781 this.control.analyze_code_impact(request, changes).await
782 })
783 }
784
785 fn analyze_code_impact_scope(
786 &self,
787 source_scope: String,
788 request: crate::domain::CodeImpactRequest,
789 changes: CodeImpactChanges,
790 ) -> StorageFuture<'_, Vec<CodeRetrievalHit>> {
791 routing::analyze_code_impact_scope(self.clone(), source_scope, request, changes)
792 }
793
794 fn codebase_view_snapshot(
795 &self,
796 source_scope: String,
797 request: crate::domain::CodebaseViewRequest,
798 row_limit: usize,
799 ) -> StorageFuture<'_, crate::domain::CodebaseViewSnapshot> {
800 routing::codebase_view_snapshot(self.clone(), source_scope, request, row_limit)
801 }
802
803 fn code_repository_totals(&self) -> StorageFuture<'_, CodeRepositoryTotals> {
804 let this = self.clone();
805 Box::pin(async move { totals::code_repository_totals(this.control, this.catalog).await })
806 }
807
808 fn code_repository_report(
809 &self,
810 repository: String,
811 ) -> StorageFuture<'_, CodeRepositoryReport> {
812 let this = self.clone();
813 Box::pin(async move {
814 if let Some(shard) =
815 repository_store_for_selector(&this.control, &this.catalog, repository.clone())
816 .await?
817 {
818 return shard.code_repository_report(repository).await;
819 }
820 this.control.code_repository_report(repository).await
821 })
822 }
823
824 fn code_repository_scope_symbol_generation_counts(
825 &self,
826 source_scope: String,
827 ) -> StorageFuture<'_, CodeSymbolGenerationCounts> {
828 let this = self.clone();
829 Box::pin(async move {
830 totals::scope_symbol_generation_counts(this.control, this.catalog, source_scope).await
831 })
832 }
833
834 fn refresh_software_global_projection(
835 &self,
836 source_scope: String,
837 ) -> StorageFuture<'_, SoftwareGlobalProjection> {
838 let this = self.clone();
839 Box::pin(async move {
840 if let Some(shard) = source_scope_store(&this.catalog, source_scope.clone()).await? {
841 return shard.refresh_software_global_projection(source_scope).await;
842 }
843 this.control
844 .refresh_software_global_projection(source_scope)
845 .await
846 })
847 }
848
849 fn software_global_projection(
850 &self,
851 request: SoftwareGlobalRequest,
852 ) -> StorageFuture<'_, SoftwareGlobalProjection> {
853 let this = self.clone();
854 Box::pin(async move {
855 if let Some(shard) = repository_store_for_selector(
856 &this.control,
857 &this.catalog,
858 request.repository.repository.clone(),
859 )
860 .await?
861 {
862 return match shard.software_global_projection(request.clone()).await {
863 Ok(projection) => Ok(projection),
864 Err(error) if is_missing_code_scope_error(&error) => {
865 this.control.software_global_projection(request).await
866 }
867 Err(error) => Err(error),
868 };
869 }
870 this.control.software_global_projection(request).await
871 })
872 }
873
874 fn software_global_projection_for_scope(
875 &self,
876 source_scope: String,
877 request: SoftwareGlobalRequest,
878 ) -> StorageFuture<'_, SoftwareGlobalProjection> {
879 let this = self.clone();
880 Box::pin(async move {
881 if let Some(shard) = source_scope_store(&this.catalog, source_scope.clone()).await? {
882 return shard
883 .software_global_projection_for_scope(source_scope, request)
884 .await;
885 }
886 this.control
887 .software_global_projection_for_scope(source_scope, request)
888 .await
889 })
890 }
891
892 fn create_code_repository_set(
893 &self,
894 seed: CodeRepositorySetSeed,
895 ) -> StorageFuture<'_, CodeRepositorySet> {
896 self.control.create_code_repository_set(seed)
897 }
898
899 fn add_code_repository_set_member(
900 &self,
901 seed: CodeRepositorySetMemberSeed,
902 ) -> StorageFuture<'_, CodeRepositorySetMember> {
903 self.control.add_code_repository_set_member(seed)
904 }
905
906 fn remove_code_repository_set_member(
907 &self,
908 set_alias: String,
909 repository_alias: String,
910 ) -> StorageFuture<'_, CodeRepositorySetMember> {
911 self.control
912 .remove_code_repository_set_member(set_alias, repository_alias)
913 }
914
915 fn code_repository_set(
916 &self,
917 set_alias: String,
918 ) -> StorageFuture<'_, Option<CodeRepositorySet>> {
919 self.control.code_repository_set(set_alias)
920 }
921
922 fn code_repository_set_status(
923 &self,
924 set_alias: String,
925 ) -> StorageFuture<'_, Option<CodeRepositorySetStatus>> {
926 self.control.code_repository_set_status(set_alias)
927 }
928
929 fn refresh_code_repository_set_overlay(
930 &self,
931 set_alias: String,
932 _now_ms: u64,
933 ) -> StorageFuture<'_, CodeRepositorySetRefreshSummary> {
934 Box::pin(async move {
935 Err(StorageError::InvalidInput(format!(
936 "repository set overlay refresh for '{set_alias}' requires the single_sqlite topology until cross-shard import/export aggregation is implemented"
937 )))
938 })
939 }
940
941 fn code_repository_set_cross_edges(
942 &self,
943 set_id: String,
944 ) -> StorageFuture<'_, Vec<CodeRepositoryCrossEdge>> {
945 self.control.code_repository_set_cross_edges(set_id)
946 }
947
948 fn queue_code_repository_set_refresh_task(
949 &self,
950 task: CodeRepositorySetRefreshTaskSeed,
951 ) -> StorageFuture<'_, crate::domain::CodeRepositorySetRefreshTaskRecord> {
952 Box::pin(async move {
953 Err(StorageError::InvalidInput(format!(
954 "repository set overlay refresh task for '{}' requires the single_sqlite topology until cross-shard import/export aggregation is implemented",
955 task.set_alias
956 )))
957 })
958 }
959
960 fn claim_code_repository_set_refresh_task(
961 &self,
962 request: CodeRepositorySetRefreshTaskClaimRequest,
963 ) -> StorageFuture<'_, Option<crate::domain::CodeRepositorySetRefreshTaskRecord>> {
964 self.control.claim_code_repository_set_refresh_task(request)
965 }
966
967 fn complete_code_repository_set_refresh_task(
968 &self,
969 request: CodeRepositorySetRefreshTaskCompletion,
970 ) -> StorageFuture<'_, crate::domain::CodeRepositorySetRefreshTaskRecord> {
971 self.control
972 .complete_code_repository_set_refresh_task(request)
973 }
974
975 fn fail_code_repository_set_refresh_task(
976 &self,
977 request: CodeRepositorySetRefreshTaskFailure,
978 ) -> StorageFuture<'_, crate::domain::CodeRepositorySetRefreshTaskRecord> {
979 self.control.fail_code_repository_set_refresh_task(request)
980 }
981}