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#[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}