Skip to main content

relay_knowledge/storage/
partitioned.rs

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, SoftwareGlobalProjection, SoftwareGlobalRequest,
26    },
27    paths::RuntimePaths,
28    storage::{
29        CodeImpactChanges, CodeIndexTaskClaimRequest, CodeIndexTaskCompletion,
30        CodeIndexTaskFailure, CodeIndexTaskLeaseRecord, CodeIndexTaskLeaseRecovery,
31        CodeIndexTaskLeaseRenewal, CodeRepositorySetMemberSeed,
32        CodeRepositorySetRefreshTaskClaimRequest, CodeRepositorySetRefreshTaskCompletion,
33        CodeRepositorySetRefreshTaskFailure, CodeRepositorySetRefreshTaskSeed,
34        CodeRepositorySetSeed, CodeRepositoryStore, CodeScopeRetentionRequest, SqliteGraphStore,
35        StorageError, StorageFuture,
36    },
37};
38
39use catalog::{SqliteShardCatalog, initialize_catalog_schema};
40use retention::merge_scope_retention_summaries;
41use routing::{
42    current_control_scope, is_missing_code_scope_error, repository_store_for_selector,
43    source_scope_store,
44};
45use status::mirror_status;
46use totals::add_code_repository_totals;
47
48/// SQLite topology that keeps global control state in one DB and code facts in
49/// one DB per registered repository.
50#[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    async fn incremental_base_scope(
69        &self,
70        snapshot: &CodeIndexSnapshot,
71    ) -> Result<Option<String>, StorageError> {
72        if snapshot.full_replace {
73            return Ok(None);
74        }
75        let Some(base_commit) = snapshot.base_resolved_commit_sha.clone() else {
76            return Ok(None);
77        };
78
79        Ok(self
80            .control
81            .code_repository_scope_status(
82                snapshot.repository_id.clone(),
83                base_commit,
84                snapshot.path_filters.clone(),
85                snapshot.language_filters.clone(),
86            )
87            .await?
88            .and_then(|status| status.last_indexed_scope_id))
89    }
90}
91
92impl CodeRepositoryStore for PartitionedSqliteKnowledgeStore {
93    fn upsert_code_repository(
94        &self,
95        registration: CodeRepositoryRegistration,
96    ) -> StorageFuture<'_, CodeRepositoryStatus> {
97        let this = self.clone();
98        Box::pin(async move {
99            let status = this
100                .control
101                .upsert_code_repository(registration.clone())
102                .await?;
103            let imported_scope = status.last_indexed_scope_id.clone();
104            let shard = this
105                .catalog
106                .staged_repository_store(status.repository_id.clone())
107                .await?;
108            this.catalog
109                .import_control_repository(
110                    Arc::clone(&shard),
111                    status.repository_id.clone(),
112                    imported_scope.clone(),
113                )
114                .await?;
115            let shard_status = shard.upsert_code_repository(registration).await?;
116            if let Some(source_scope) = imported_scope {
117                this.catalog
118                    .record_scope(status.repository_id.clone(), source_scope)
119                    .await?;
120            } else {
121                this.catalog
122                    .activate_repository(status.repository_id.clone())
123                    .await?;
124            }
125            Ok(CodeRepositoryStatus {
126                alias: status.alias,
127                ..shard_status
128            })
129        })
130    }
131
132    fn code_repository_status(
133        &self,
134        repository: String,
135    ) -> StorageFuture<'_, Option<CodeRepositoryStatus>> {
136        let this = self.clone();
137        Box::pin(async move {
138            let Some(control_status) = this.control.code_repository_status(repository).await?
139            else {
140                return Ok(None);
141            };
142            let Some(shard) = this
143                .catalog
144                .existing_repository_store(control_status.repository_id.clone())
145                .await?
146            else {
147                return Ok(Some(control_status));
148            };
149            let Some(mut shard_status) = shard
150                .code_repository_status(control_status.repository_id.clone())
151                .await?
152            else {
153                return Ok(Some(control_status));
154            };
155            shard_status.alias = control_status.alias;
156            Ok(Some(shard_status))
157        })
158    }
159
160    fn remove_code_repository(
161        &self,
162        repository: String,
163        now_ms: u64,
164    ) -> StorageFuture<'_, Option<CodeRepositoryRemovalSummary>> {
165        let this = self.clone();
166        Box::pin(async move {
167            let Some(control_status) = this.control.code_repository_status(repository).await?
168            else {
169                return Ok(None);
170            };
171            let shard = this
172                .catalog
173                .existing_repository_store(control_status.repository_id.clone())
174                .await?;
175            let removed = this
176                .control
177                .remove_code_repository(control_status.repository_id.clone(), now_ms)
178                .await?;
179            let Some(summary) = removed else {
180                return Ok(None);
181            };
182            if let Some(shard) = shard {
183                shard
184                    .remove_code_repository(control_status.repository_id.clone(), now_ms)
185                    .await?;
186            }
187            this.catalog
188                .remove_repository(control_status.repository_id)
189                .await?;
190            Ok(Some(summary))
191        })
192    }
193
194    fn code_repository_scope_status(
195        &self,
196        repository: String,
197        resolved_commit_sha: String,
198        path_filters: Vec<String>,
199        language_filters: Vec<String>,
200    ) -> StorageFuture<'_, Option<CodeRepositoryStatus>> {
201        let this = self.clone();
202        Box::pin(async move {
203            let Some(control_status) = this.control.code_repository_status(repository).await?
204            else {
205                return Ok(None);
206            };
207            let Some(shard) = this
208                .catalog
209                .existing_repository_store(control_status.repository_id.clone())
210                .await?
211            else {
212                return this
213                    .control
214                    .code_repository_scope_status(
215                        control_status.repository_id,
216                        resolved_commit_sha,
217                        path_filters,
218                        language_filters,
219                    )
220                    .await;
221            };
222            let status = shard
223                .code_repository_scope_status(
224                    control_status.repository_id.clone(),
225                    resolved_commit_sha.clone(),
226                    path_filters.clone(),
227                    language_filters.clone(),
228                )
229                .await?;
230            if let Some(mut status) = status {
231                status.alias = control_status.alias;
232                return Ok(Some(status));
233            }
234            this.control
235                .code_repository_scope_status(
236                    control_status.repository_id,
237                    resolved_commit_sha,
238                    path_filters,
239                    language_filters,
240                )
241                .await
242        })
243    }
244
245    fn latest_code_repository_scope_status(
246        &self,
247        repository: String,
248        path_filters: Vec<String>,
249        language_filters: Vec<String>,
250    ) -> StorageFuture<'_, Option<CodeRepositoryStatus>> {
251        let this = self.clone();
252        Box::pin(async move {
253            let Some(control_status) = this.control.code_repository_status(repository).await?
254            else {
255                return Ok(None);
256            };
257            let Some(shard) = this
258                .catalog
259                .existing_repository_store(control_status.repository_id.clone())
260                .await?
261            else {
262                return this
263                    .control
264                    .latest_code_repository_scope_status(
265                        control_status.repository_id,
266                        path_filters,
267                        language_filters,
268                    )
269                    .await;
270            };
271            let status = shard
272                .latest_code_repository_scope_status(
273                    control_status.repository_id.clone(),
274                    path_filters.clone(),
275                    language_filters.clone(),
276                )
277                .await?;
278            if let Some(mut status) = status {
279                status.alias = control_status.alias;
280                return Ok(Some(status));
281            }
282            this.control
283                .latest_code_repository_scope_status(
284                    control_status.repository_id,
285                    path_filters,
286                    language_filters,
287                )
288                .await
289        })
290    }
291
292    fn queue_code_index_task(
293        &self,
294        task: crate::storage::CodeIndexTaskSeed,
295    ) -> StorageFuture<'_, crate::domain::CodeIndexTaskRecord> {
296        let control = Arc::clone(&self.control);
297        Box::pin(async move { control.queue_code_index_task(task).await })
298    }
299
300    fn claim_code_index_task(
301        &self,
302        request: CodeIndexTaskClaimRequest,
303    ) -> StorageFuture<'_, Option<crate::domain::CodeIndexTaskRecord>> {
304        self.control.claim_code_index_task(request)
305    }
306
307    fn recover_code_index_task_leases(
308        &self,
309        now_ms: u64,
310        max_attempts: u32,
311    ) -> StorageFuture<'_, ()> {
312        self.control
313            .recover_code_index_task_leases(now_ms, max_attempts)
314    }
315
316    fn running_code_index_task_leases(&self) -> StorageFuture<'_, Vec<CodeIndexTaskLeaseRecord>> {
317        self.control.running_code_index_task_leases()
318    }
319
320    fn recover_code_index_task_leases_by_task(
321        &self,
322        request: CodeIndexTaskLeaseRecovery,
323    ) -> StorageFuture<'_, usize> {
324        self.control.recover_code_index_task_leases_by_task(request)
325    }
326
327    fn reset_code_index_tasks(
328        &self,
329        repository_id: String,
330        now_ms: u64,
331    ) -> StorageFuture<'_, Vec<crate::domain::CodeIndexTaskRecord>> {
332        self.control.reset_code_index_tasks(repository_id, now_ms)
333    }
334
335    fn renew_code_index_task_lease(
336        &self,
337        request: CodeIndexTaskLeaseRenewal,
338    ) -> StorageFuture<'_, crate::domain::CodeIndexTaskRecord> {
339        self.control.renew_code_index_task_lease(request)
340    }
341
342    fn complete_code_index_task(
343        &self,
344        request: CodeIndexTaskCompletion,
345    ) -> StorageFuture<'_, crate::domain::CodeIndexTaskRecord> {
346        self.control.complete_code_index_task(request)
347    }
348
349    fn fail_code_index_task(
350        &self,
351        request: CodeIndexTaskFailure,
352    ) -> StorageFuture<'_, crate::domain::CodeIndexTaskRecord> {
353        self.control.fail_code_index_task(request)
354    }
355
356    fn code_index_task(
357        &self,
358        task_id: String,
359    ) -> StorageFuture<'_, Option<crate::domain::CodeIndexTaskRecord>> {
360        self.control.code_index_task(task_id)
361    }
362
363    fn active_code_index_task(
364        &self,
365        repository_id: String,
366    ) -> StorageFuture<'_, Option<crate::domain::CodeIndexTaskRecord>> {
367        self.control.active_code_index_task(repository_id)
368    }
369
370    fn code_index_task_queue_status(
371        &self,
372    ) -> StorageFuture<'_, crate::domain::CodeIndexTaskQueueStatus> {
373        self.control.code_index_task_queue_status()
374    }
375    fn code_index_checkpoint(
376        &self,
377        source_scope: 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_scope_store(source_scope.clone())
384                .await?
385            {
386                if let Some(checkpoint) = shard.code_index_checkpoint(source_scope.clone()).await? {
387                    return Ok(Some(checkpoint));
388                }
389            }
390            this.control.code_index_checkpoint(source_scope).await
391        })
392    }
393
394    fn latest_code_index_checkpoint(
395        &self,
396        repository_id: String,
397    ) -> StorageFuture<'_, Option<CodeIndexCheckpoint>> {
398        let this = self.clone();
399        Box::pin(async move {
400            if let Some(shard) = this
401                .catalog
402                .checkpoint_repository_store(repository_id.clone())
403                .await?
404            {
405                return shard.latest_code_index_checkpoint(repository_id).await;
406            }
407            this.control
408                .latest_code_index_checkpoint(repository_id)
409                .await
410        })
411    }
412
413    fn code_scope_retention(
414        &self,
415        repository_id: String,
416    ) -> StorageFuture<'_, crate::domain::CodeScopeRetentionSummary> {
417        let this = self.clone();
418        Box::pin(async move {
419            if let Some(shard) = this
420                .catalog
421                .existing_repository_store(repository_id.clone())
422                .await?
423            {
424                return shard.code_scope_retention(repository_id).await;
425            }
426            this.control.code_scope_retention(repository_id).await
427        })
428    }
429
430    fn prune_code_repository_scopes(
431        &self,
432        request: CodeScopeRetentionRequest,
433    ) -> StorageFuture<'_, crate::domain::CodeScopeRetentionSummary> {
434        let this = self.clone();
435        Box::pin(async move {
436            if let Some(shard) = this
437                .catalog
438                .existing_repository_store(request.repository_id.clone())
439                .await?
440            {
441                let control_retention = this
442                    .control
443                    .prune_code_repository_scopes(request.clone())
444                    .await?;
445                let shard_retention = shard
446                    .prune_code_repository_scopes_with_retained(
447                        request.clone(),
448                        control_retention.retained_scopes.clone(),
449                    )
450                    .await;
451                return shard_retention.map(|summary| {
452                    merge_scope_retention_summaries(
453                        request.repository_id,
454                        control_retention,
455                        summary,
456                    )
457                });
458            }
459            this.control.prune_code_repository_scopes(request).await
460        })
461    }
462
463    fn code_file_fingerprints(
464        &self,
465        repository_id: String,
466    ) -> StorageFuture<'_, Vec<crate::domain::CodeFileFingerprint>> {
467        let this = self.clone();
468        Box::pin(async move {
469            if let Some(shard) = this
470                .catalog
471                .existing_repository_store(repository_id.clone())
472                .await?
473            {
474                return shard.code_file_fingerprints(repository_id).await;
475            }
476            this.control.code_file_fingerprints(repository_id).await
477        })
478    }
479
480    fn code_file_fingerprints_for_scope(
481        &self,
482        source_scope: String,
483    ) -> StorageFuture<'_, Vec<crate::domain::CodeFileFingerprint>> {
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.code_file_fingerprints_for_scope(source_scope).await;
488            }
489            this.control
490                .code_file_fingerprints_for_scope(source_scope)
491                .await
492        })
493    }
494
495    fn code_file_candidate_paths_for_scope(
496        &self,
497        source_scope: String,
498        path_filters: Vec<String>,
499        language_filters: Vec<String>,
500        limit: usize,
501    ) -> StorageFuture<'_, Vec<String>> {
502        let this = self.clone();
503        Box::pin(async move {
504            if let Some(shard) = source_scope_store(&this.catalog, source_scope.clone()).await? {
505                return shard
506                    .code_file_candidate_paths_for_scope(
507                        source_scope,
508                        path_filters,
509                        language_filters,
510                        limit,
511                    )
512                    .await;
513            }
514            this.control
515                .code_file_candidate_paths_for_scope(
516                    source_scope,
517                    path_filters,
518                    language_filters,
519                    limit,
520                )
521                .await
522        })
523    }
524
525    fn code_file_candidate_paths_for_query_scope(
526        &self,
527        source_scope: String,
528        query: String,
529        path_filters: Vec<String>,
530        language_filters: Vec<String>,
531        limit: usize,
532    ) -> StorageFuture<'_, Vec<String>> {
533        let this = self.clone();
534        Box::pin(async move {
535            if let Some(shard) = source_scope_store(&this.catalog, source_scope.clone()).await? {
536                return shard
537                    .code_file_candidate_paths_for_query_scope(
538                        source_scope,
539                        query,
540                        path_filters,
541                        language_filters,
542                        limit,
543                    )
544                    .await;
545            }
546            this.control
547                .code_file_candidate_paths_for_query_scope(
548                    source_scope,
549                    query,
550                    path_filters,
551                    language_filters,
552                    limit,
553                )
554                .await
555        })
556    }
557
558    fn apply_code_index_snapshot(
559        &self,
560        snapshot: CodeIndexSnapshot,
561    ) -> StorageFuture<'_, CodeIndexSummary> {
562        let this = self.clone();
563        Box::pin(async move {
564            let base_scope = this.incremental_base_scope(&snapshot).await?;
565            let shard = if snapshot.full_replace {
566                this.catalog
567                    .staged_repository_store(snapshot.repository_id.clone())
568                    .await?
569            } else {
570                match this
571                    .catalog
572                    .existing_repository_store(snapshot.repository_id.clone())
573                    .await?
574                {
575                    Some(shard) => shard,
576                    None => {
577                        this.catalog
578                            .staged_repository_store(snapshot.repository_id.clone())
579                            .await?
580                    }
581                }
582            };
583            this.catalog
584                .import_control_repository(
585                    Arc::clone(&shard),
586                    snapshot.repository_id.clone(),
587                    base_scope,
588                )
589                .await?;
590            let summary = shard.apply_code_index_snapshot(snapshot).await?;
591            let status = shard
592                .code_repository_status(summary.repository_id.clone())
593                .await?
594                .ok_or_else(|| {
595                    StorageError::InvalidInput(
596                        "sharded code repository status is missing after index".to_owned(),
597                    )
598                })?;
599            this.catalog
600                .record_scope(summary.repository_id.clone(), summary.source_scope.clone())
601                .await?;
602            mirror_status(&this.control, status).await?;
603            Ok(summary)
604        })
605    }
606
607    fn begin_code_index_session(
608        &self,
609        session: CodeIndexSession,
610    ) -> StorageFuture<'_, CodeIndexCheckpoint> {
611        let this = self.clone();
612        Box::pin(async move {
613            let repository_id = session.repository_id.clone();
614            let source_scope = session.source_scope.clone();
615            let shard = this
616                .catalog
617                .staged_repository_store(repository_id.clone())
618                .await?;
619            let control_scope = current_control_scope(&this.control, repository_id.clone()).await?;
620            this.catalog
621                .import_control_repository(Arc::clone(&shard), repository_id.clone(), control_scope)
622                .await?;
623            let checkpoint = shard.begin_code_index_session(session).await?;
624            this.catalog
625                .stage_scope(repository_id, source_scope)
626                .await?;
627            Ok(checkpoint)
628        })
629    }
630
631    fn apply_code_index_batch(
632        &self,
633        batch: CodeIndexBatch,
634    ) -> StorageFuture<'_, CodeIndexCheckpoint> {
635        let this = self.clone();
636        Box::pin(async move {
637            let repository_id = batch.repository_id.clone();
638            let source_scope = batch.source_scope.clone();
639            let shard = this
640                .catalog
641                .staged_repository_store(repository_id.clone())
642                .await?;
643            let control_scope = current_control_scope(&this.control, repository_id.clone()).await?;
644            this.catalog
645                .import_control_repository(Arc::clone(&shard), repository_id.clone(), control_scope)
646                .await?;
647            let checkpoint = shard.apply_code_index_batch(batch).await?;
648            this.catalog
649                .stage_scope(repository_id, source_scope)
650                .await?;
651            Ok(checkpoint)
652        })
653    }
654
655    fn finalize_code_index_session(
656        &self,
657        session: CodeIndexSession,
658    ) -> StorageFuture<'_, CodeIndexSummary> {
659        let this = self.clone();
660        Box::pin(async move {
661            let shard = this
662                .catalog
663                .staged_repository_store(session.repository_id.clone())
664                .await?;
665            let summary = shard.finalize_code_index_session(session).await?;
666            let status = shard
667                .code_repository_status(summary.repository_id.clone())
668                .await?
669                .ok_or_else(|| {
670                    StorageError::InvalidInput(
671                        "sharded code repository status is missing after finalize".to_owned(),
672                    )
673                })?;
674            this.catalog
675                .record_scope(summary.repository_id.clone(), summary.source_scope.clone())
676                .await?;
677            mirror_status(&this.control, status).await?;
678            Ok(summary)
679        })
680    }
681
682    fn search_code(
683        &self,
684        request: CodeRetrievalRequest,
685    ) -> StorageFuture<'_, Vec<CodeRetrievalHit>> {
686        let this = self.clone();
687        Box::pin(async move {
688            if let Some(shard) = repository_store_for_selector(
689                &this.control,
690                &this.catalog,
691                request.repository.repository.clone(),
692            )
693            .await?
694            {
695                return match shard.search_code(request.clone()).await {
696                    Ok(hits) => Ok(hits),
697                    Err(error) if is_missing_code_scope_error(&error) => {
698                        this.control.search_code(request).await
699                    }
700                    Err(error) => Err(error),
701                };
702            }
703            this.control.search_code(request).await
704        })
705    }
706
707    fn search_code_feature_flags(
708        &self,
709        request: CodeFeatureFlagRequest,
710    ) -> StorageFuture<'_, Vec<CodeFeatureFlagGraph>> {
711        let this = self.clone();
712        Box::pin(async move {
713            if let Some(shard) = repository_store_for_selector(
714                &this.control,
715                &this.catalog,
716                request.repository.repository.clone(),
717            )
718            .await?
719            {
720                return match shard.search_code_feature_flags(request.clone()).await {
721                    Ok(flags) => Ok(flags),
722                    Err(error) if is_missing_code_scope_error(&error) => {
723                        this.control.search_code_feature_flags(request).await
724                    }
725                    Err(error) => Err(error),
726                };
727            }
728            this.control.search_code_feature_flags(request).await
729        })
730    }
731
732    fn search_code_feature_flags_scope(
733        &self,
734        source_scope: String,
735        request: CodeFeatureFlagRequest,
736    ) -> StorageFuture<'_, Vec<CodeFeatureFlagGraph>> {
737        let this = self.clone();
738        Box::pin(async move {
739            if let Some(shard) = source_scope_store(&this.catalog, source_scope.clone()).await? {
740                return shard
741                    .search_code_feature_flags_scope(source_scope, request)
742                    .await;
743            }
744            this.control
745                .search_code_feature_flags_scope(source_scope, request)
746                .await
747        })
748    }
749
750    fn search_code_scope(
751        &self,
752        source_scope: String,
753        request: CodeRetrievalRequest,
754    ) -> StorageFuture<'_, Vec<CodeRetrievalHit>> {
755        let this = self.clone();
756        Box::pin(async move {
757            if let Some(shard) = source_scope_store(&this.catalog, source_scope.clone()).await? {
758                return shard.search_code_scope(source_scope, request).await;
759            }
760            this.control.search_code_scope(source_scope, request).await
761        })
762    }
763
764    fn analyze_code_impact(
765        &self,
766        request: crate::domain::CodeImpactRequest,
767        changes: CodeImpactChanges,
768    ) -> StorageFuture<'_, Vec<CodeRetrievalHit>> {
769        let this = self.clone();
770        Box::pin(async move {
771            if let Some(shard) = repository_store_for_selector(
772                &this.control,
773                &this.catalog,
774                request.repository.repository.clone(),
775            )
776            .await?
777            {
778                return match shard
779                    .analyze_code_impact(request.clone(), changes.clone())
780                    .await
781                {
782                    Ok(hits) => Ok(hits),
783                    Err(error) if is_missing_code_scope_error(&error) => {
784                        this.control.analyze_code_impact(request, changes).await
785                    }
786                    Err(error) => Err(error),
787                };
788            }
789            this.control.analyze_code_impact(request, changes).await
790        })
791    }
792
793    fn analyze_code_impact_scope(
794        &self,
795        source_scope: String,
796        request: crate::domain::CodeImpactRequest,
797        changes: CodeImpactChanges,
798    ) -> StorageFuture<'_, Vec<CodeRetrievalHit>> {
799        let this = self.clone();
800        Box::pin(async move {
801            if let Some(shard) = source_scope_store(&this.catalog, source_scope.clone()).await? {
802                return shard
803                    .analyze_code_impact_scope(source_scope, request, changes)
804                    .await;
805            }
806            this.control
807                .analyze_code_impact_scope(source_scope, request, changes)
808                .await
809        })
810    }
811
812    fn code_repository_totals(&self) -> StorageFuture<'_, CodeRepositoryTotals> {
813        let this = self.clone();
814        Box::pin(async move {
815            let repository_ids = this.catalog.repository_ids().await?;
816            let mut totals = this
817                .control
818                .code_repository_totals_excluding(repository_ids.clone())
819                .await?;
820            for repository_id in repository_ids {
821                let Some(shard) = this
822                    .catalog
823                    .existing_repository_store(repository_id)
824                    .await?
825                else {
826                    continue;
827                };
828                let shard_totals = shard.code_repository_totals().await?;
829                add_code_repository_totals(&mut totals, shard_totals);
830            }
831            Ok(totals)
832        })
833    }
834
835    fn code_repository_report(
836        &self,
837        repository: String,
838    ) -> StorageFuture<'_, CodeRepositoryReport> {
839        let this = self.clone();
840        Box::pin(async move {
841            if let Some(shard) =
842                repository_store_for_selector(&this.control, &this.catalog, repository.clone())
843                    .await?
844            {
845                return shard.code_repository_report(repository).await;
846            }
847            this.control.code_repository_report(repository).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}