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 let this = self.clone();
746 Box::pin(async move {
747 if let Some(shard) = source_scope_store(&this.catalog, source_scope.clone()).await? {
748 return shard
749 .search_code_feature_flags_scope(source_scope, request)
750 .await;
751 }
752 this.control
753 .search_code_feature_flags_scope(source_scope, request)
754 .await
755 })
756 }
757
758 fn search_code_scope(
759 &self,
760 source_scope: String,
761 request: CodeRetrievalRequest,
762 ) -> StorageFuture<'_, Vec<CodeRetrievalHit>> {
763 let this = self.clone();
764 Box::pin(async move {
765 if let Some(shard) = source_scope_store(&this.catalog, source_scope.clone()).await? {
766 return shard.search_code_scope(source_scope, request).await;
767 }
768 this.control.search_code_scope(source_scope, request).await
769 })
770 }
771
772 fn analyze_code_impact(
773 &self,
774 request: crate::domain::CodeImpactRequest,
775 changes: CodeImpactChanges,
776 ) -> StorageFuture<'_, Vec<CodeRetrievalHit>> {
777 let this = self.clone();
778 Box::pin(async move {
779 if let Some(shard) = repository_store_for_selector(
780 &this.control,
781 &this.catalog,
782 request.repository.repository.clone(),
783 )
784 .await?
785 {
786 return match shard
787 .analyze_code_impact(request.clone(), changes.clone())
788 .await
789 {
790 Ok(hits) => Ok(hits),
791 Err(error) if is_missing_code_scope_error(&error) => {
792 this.control.analyze_code_impact(request, changes).await
793 }
794 Err(error) => Err(error),
795 };
796 }
797 this.control.analyze_code_impact(request, changes).await
798 })
799 }
800
801 fn analyze_code_impact_scope(
802 &self,
803 source_scope: String,
804 request: crate::domain::CodeImpactRequest,
805 changes: CodeImpactChanges,
806 ) -> StorageFuture<'_, Vec<CodeRetrievalHit>> {
807 let this = self.clone();
808 Box::pin(async move {
809 if let Some(shard) = source_scope_store(&this.catalog, source_scope.clone()).await? {
810 return shard
811 .analyze_code_impact_scope(source_scope, request, changes)
812 .await;
813 }
814 this.control
815 .analyze_code_impact_scope(source_scope, request, changes)
816 .await
817 })
818 }
819
820 fn code_repository_totals(&self) -> StorageFuture<'_, CodeRepositoryTotals> {
821 let this = self.clone();
822 Box::pin(async move { totals::code_repository_totals(this.control, this.catalog).await })
823 }
824
825 fn code_repository_report(
826 &self,
827 repository: String,
828 ) -> StorageFuture<'_, CodeRepositoryReport> {
829 let this = self.clone();
830 Box::pin(async move {
831 if let Some(shard) =
832 repository_store_for_selector(&this.control, &this.catalog, repository.clone())
833 .await?
834 {
835 return shard.code_repository_report(repository).await;
836 }
837 this.control.code_repository_report(repository).await
838 })
839 }
840
841 fn code_repository_scope_symbol_generation_counts(
842 &self,
843 source_scope: String,
844 ) -> StorageFuture<'_, CodeSymbolGenerationCounts> {
845 let this = self.clone();
846 Box::pin(async move {
847 totals::scope_symbol_generation_counts(this.control, this.catalog, source_scope).await
848 })
849 }
850
851 fn refresh_software_global_projection(
852 &self,
853 source_scope: String,
854 ) -> StorageFuture<'_, SoftwareGlobalProjection> {
855 let this = self.clone();
856 Box::pin(async move {
857 if let Some(shard) = source_scope_store(&this.catalog, source_scope.clone()).await? {
858 return shard.refresh_software_global_projection(source_scope).await;
859 }
860 this.control
861 .refresh_software_global_projection(source_scope)
862 .await
863 })
864 }
865
866 fn software_global_projection(
867 &self,
868 request: SoftwareGlobalRequest,
869 ) -> StorageFuture<'_, SoftwareGlobalProjection> {
870 let this = self.clone();
871 Box::pin(async move {
872 if let Some(shard) = repository_store_for_selector(
873 &this.control,
874 &this.catalog,
875 request.repository.repository.clone(),
876 )
877 .await?
878 {
879 return match shard.software_global_projection(request.clone()).await {
880 Ok(projection) => Ok(projection),
881 Err(error) if is_missing_code_scope_error(&error) => {
882 this.control.software_global_projection(request).await
883 }
884 Err(error) => Err(error),
885 };
886 }
887 this.control.software_global_projection(request).await
888 })
889 }
890
891 fn software_global_projection_for_scope(
892 &self,
893 source_scope: String,
894 request: SoftwareGlobalRequest,
895 ) -> StorageFuture<'_, SoftwareGlobalProjection> {
896 let this = self.clone();
897 Box::pin(async move {
898 if let Some(shard) = source_scope_store(&this.catalog, source_scope.clone()).await? {
899 return shard
900 .software_global_projection_for_scope(source_scope, request)
901 .await;
902 }
903 this.control
904 .software_global_projection_for_scope(source_scope, request)
905 .await
906 })
907 }
908
909 fn create_code_repository_set(
910 &self,
911 seed: CodeRepositorySetSeed,
912 ) -> StorageFuture<'_, CodeRepositorySet> {
913 self.control.create_code_repository_set(seed)
914 }
915
916 fn add_code_repository_set_member(
917 &self,
918 seed: CodeRepositorySetMemberSeed,
919 ) -> StorageFuture<'_, CodeRepositorySetMember> {
920 self.control.add_code_repository_set_member(seed)
921 }
922
923 fn remove_code_repository_set_member(
924 &self,
925 set_alias: String,
926 repository_alias: String,
927 ) -> StorageFuture<'_, CodeRepositorySetMember> {
928 self.control
929 .remove_code_repository_set_member(set_alias, repository_alias)
930 }
931
932 fn code_repository_set(
933 &self,
934 set_alias: String,
935 ) -> StorageFuture<'_, Option<CodeRepositorySet>> {
936 self.control.code_repository_set(set_alias)
937 }
938
939 fn code_repository_set_status(
940 &self,
941 set_alias: String,
942 ) -> StorageFuture<'_, Option<CodeRepositorySetStatus>> {
943 self.control.code_repository_set_status(set_alias)
944 }
945
946 fn refresh_code_repository_set_overlay(
947 &self,
948 set_alias: String,
949 _now_ms: u64,
950 ) -> StorageFuture<'_, CodeRepositorySetRefreshSummary> {
951 Box::pin(async move {
952 Err(StorageError::InvalidInput(format!(
953 "repository set overlay refresh for '{set_alias}' requires the single_sqlite topology until cross-shard import/export aggregation is implemented"
954 )))
955 })
956 }
957
958 fn code_repository_set_cross_edges(
959 &self,
960 set_id: String,
961 ) -> StorageFuture<'_, Vec<CodeRepositoryCrossEdge>> {
962 self.control.code_repository_set_cross_edges(set_id)
963 }
964
965 fn queue_code_repository_set_refresh_task(
966 &self,
967 task: CodeRepositorySetRefreshTaskSeed,
968 ) -> StorageFuture<'_, crate::domain::CodeRepositorySetRefreshTaskRecord> {
969 Box::pin(async move {
970 Err(StorageError::InvalidInput(format!(
971 "repository set overlay refresh task for '{}' requires the single_sqlite topology until cross-shard import/export aggregation is implemented",
972 task.set_alias
973 )))
974 })
975 }
976
977 fn claim_code_repository_set_refresh_task(
978 &self,
979 request: CodeRepositorySetRefreshTaskClaimRequest,
980 ) -> StorageFuture<'_, Option<crate::domain::CodeRepositorySetRefreshTaskRecord>> {
981 self.control.claim_code_repository_set_refresh_task(request)
982 }
983
984 fn complete_code_repository_set_refresh_task(
985 &self,
986 request: CodeRepositorySetRefreshTaskCompletion,
987 ) -> StorageFuture<'_, crate::domain::CodeRepositorySetRefreshTaskRecord> {
988 self.control
989 .complete_code_repository_set_refresh_task(request)
990 }
991
992 fn fail_code_repository_set_refresh_task(
993 &self,
994 request: CodeRepositorySetRefreshTaskFailure,
995 ) -> StorageFuture<'_, crate::domain::CodeRepositorySetRefreshTaskRecord> {
996 self.control.fail_code_repository_set_refresh_task(request)
997 }
998}