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, 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/// 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
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}