Skip to main content

relay_knowledge/storage/partitioned/
mod.rs

1use std::{path::Path, sync::Arc};
2
3mod catalog;
4mod control_plane;
5mod diagnostics;
6mod indexing;
7mod repository;
8mod routing;
9mod status;
10mod totals;
11
12use crate::{
13    domain::{
14        CodeFeatureFlagGraph, CodeFeatureFlagRequest, CodeIndexBatch, CodeIndexCheckpoint,
15        CodeIndexSession, CodeIndexSnapshot, CodeIndexSummary, CodeRepositoryCrossEdge,
16        CodeRepositoryRegistration, CodeRepositoryRemovalSummary, CodeRepositoryReport,
17        CodeRepositorySet, CodeRepositorySetMember, CodeRepositorySetRefreshSummary,
18        CodeRepositorySetStatus, CodeRepositoryStatus, CodeRepositoryTotals, CodeRetrievalHit,
19        CodeRetrievalRequest, CodeSymbolGenerationCounts, SoftwareGlobalProjection,
20        SoftwareGlobalRequest,
21    },
22    paths::RuntimePaths,
23    storage::{
24        CodeImpactChanges, CodeIndexTaskClaimRequest, CodeIndexTaskCompletion,
25        CodeIndexTaskFailure, CodeIndexTaskLeaseRecord, CodeIndexTaskLeaseRecovery,
26        CodeIndexTaskLeaseRenewal, CodeRepositorySetMemberSeed,
27        CodeRepositorySetRefreshTaskClaimRequest, CodeRepositorySetRefreshTaskCompletion,
28        CodeRepositorySetRefreshTaskFailure, CodeRepositorySetRefreshTaskSeed,
29        CodeRepositorySetSeed, CodeRepositoryStore, CodeScopeRetentionRequest, SqliteGraphStore,
30        StorageError, StorageFuture,
31    },
32};
33
34use catalog::{SqliteShardCatalog, initialize_catalog_schema};
35use routing::{is_missing_code_scope_error, repository_store_for_selector, source_scope_store};
36
37/// SQLite topology that keeps global control state in one DB and code facts in
38/// one DB per registered repository.
39#[derive(Clone)]
40pub struct PartitionedSqliteKnowledgeStore {
41    control: Arc<SqliteGraphStore>,
42    catalog: Arc<SqliteShardCatalog>,
43}
44
45impl PartitionedSqliteKnowledgeStore {
46    pub fn open(control_path: impl AsRef<Path>, paths: RuntimePaths) -> Result<Self, StorageError> {
47        let control_path = control_path.as_ref().to_path_buf();
48        let control = Arc::new(SqliteGraphStore::open(&control_path)?);
49        initialize_catalog_schema(&control_path)?;
50
51        Ok(Self {
52            control,
53            catalog: Arc::new(SqliteShardCatalog::new(control_path, paths)),
54        })
55    }
56}
57
58impl CodeRepositoryStore for PartitionedSqliteKnowledgeStore {
59    fn upsert_code_repository(
60        &self,
61        registration: CodeRepositoryRegistration,
62    ) -> StorageFuture<'_, CodeRepositoryStatus> {
63        repository::upsert(self, registration)
64    }
65
66    fn code_repository_status(
67        &self,
68        repository: String,
69    ) -> StorageFuture<'_, Option<CodeRepositoryStatus>> {
70        repository::status(self, repository)
71    }
72
73    fn list_code_repositories(&self) -> StorageFuture<'_, Vec<CodeRepositoryStatus>> {
74        control_plane::list_code_repositories(self)
75    }
76
77    fn remove_code_repository(
78        &self,
79        repository: String,
80        now_ms: u64,
81    ) -> StorageFuture<'_, Option<CodeRepositoryRemovalSummary>> {
82        repository::remove(self, repository, now_ms)
83    }
84
85    fn code_repository_scope_status(
86        &self,
87        repository: String,
88        resolved_commit_sha: String,
89        path_filters: Vec<String>,
90        language_filters: Vec<String>,
91    ) -> StorageFuture<'_, Option<CodeRepositoryStatus>> {
92        repository::scope_status(
93            self,
94            repository,
95            resolved_commit_sha,
96            path_filters,
97            language_filters,
98        )
99    }
100
101    fn latest_code_repository_scope_status(
102        &self,
103        repository: String,
104        path_filters: Vec<String>,
105        language_filters: Vec<String>,
106    ) -> StorageFuture<'_, Option<CodeRepositoryStatus>> {
107        repository::latest_scope_status(self, repository, path_filters, language_filters)
108    }
109
110    fn queue_code_index_task(
111        &self,
112        task: crate::storage::CodeIndexTaskSeed,
113    ) -> StorageFuture<'_, crate::domain::CodeIndexTaskRecord> {
114        let control = Arc::clone(&self.control);
115        Box::pin(async move { control.queue_code_index_task(task).await })
116    }
117
118    fn claim_code_index_task(
119        &self,
120        request: CodeIndexTaskClaimRequest,
121    ) -> StorageFuture<'_, Option<crate::domain::CodeIndexTaskRecord>> {
122        self.control.claim_code_index_task(request)
123    }
124
125    fn recover_code_index_task_leases(
126        &self,
127        now_ms: u64,
128        max_attempts: u32,
129    ) -> StorageFuture<'_, ()> {
130        self.control
131            .recover_code_index_task_leases(now_ms, max_attempts)
132    }
133
134    fn running_code_index_task_leases(&self) -> StorageFuture<'_, Vec<CodeIndexTaskLeaseRecord>> {
135        self.control.running_code_index_task_leases()
136    }
137
138    fn recover_code_index_task_leases_by_task(
139        &self,
140        request: CodeIndexTaskLeaseRecovery,
141    ) -> StorageFuture<'_, usize> {
142        self.control.recover_code_index_task_leases_by_task(request)
143    }
144
145    fn reset_code_index_tasks(
146        &self,
147        repository_id: String,
148        now_ms: u64,
149    ) -> StorageFuture<'_, Vec<crate::domain::CodeIndexTaskRecord>> {
150        self.control.reset_code_index_tasks(repository_id, now_ms)
151    }
152
153    fn renew_code_index_task_lease(
154        &self,
155        request: CodeIndexTaskLeaseRenewal,
156    ) -> StorageFuture<'_, crate::domain::CodeIndexTaskRecord> {
157        self.control.renew_code_index_task_lease(request)
158    }
159
160    fn complete_code_index_task(
161        &self,
162        request: CodeIndexTaskCompletion,
163    ) -> StorageFuture<'_, crate::domain::CodeIndexTaskRecord> {
164        self.control.complete_code_index_task(request)
165    }
166
167    fn fail_code_index_task(
168        &self,
169        request: CodeIndexTaskFailure,
170    ) -> StorageFuture<'_, crate::domain::CodeIndexTaskRecord> {
171        self.control.fail_code_index_task(request)
172    }
173
174    fn code_index_task(
175        &self,
176        task_id: String,
177    ) -> StorageFuture<'_, Option<crate::domain::CodeIndexTaskRecord>> {
178        self.control.code_index_task(task_id)
179    }
180
181    fn active_code_index_task(
182        &self,
183        repository_id: String,
184    ) -> StorageFuture<'_, Option<crate::domain::CodeIndexTaskRecord>> {
185        self.control.active_code_index_task(repository_id)
186    }
187
188    fn code_index_task_queue_status(
189        &self,
190    ) -> StorageFuture<'_, crate::domain::CodeIndexTaskQueueStatus> {
191        self.control.code_index_task_queue_status()
192    }
193    fn code_index_checkpoint(
194        &self,
195        source_scope: String,
196    ) -> StorageFuture<'_, Option<CodeIndexCheckpoint>> {
197        indexing::checkpoint::by_scope(self, source_scope)
198    }
199
200    fn latest_code_index_checkpoint(
201        &self,
202        repository_id: String,
203    ) -> StorageFuture<'_, Option<CodeIndexCheckpoint>> {
204        indexing::checkpoint::latest(self, repository_id)
205    }
206
207    fn code_scope_retention(
208        &self,
209        repository_id: String,
210    ) -> StorageFuture<'_, crate::domain::CodeScopeRetentionSummary> {
211        indexing::retention::status(self, repository_id)
212    }
213
214    fn prune_code_repository_scopes(
215        &self,
216        request: CodeScopeRetentionRequest,
217    ) -> StorageFuture<'_, crate::domain::CodeScopeRetentionSummary> {
218        indexing::retention::prune(self, request)
219    }
220
221    fn code_file_fingerprints(
222        &self,
223        repository_id: String,
224    ) -> StorageFuture<'_, Vec<crate::domain::CodeFileFingerprint>> {
225        indexing::file_index::fingerprints(self, repository_id)
226    }
227
228    fn code_file_fingerprints_for_scope(
229        &self,
230        source_scope: String,
231    ) -> StorageFuture<'_, Vec<crate::domain::CodeFileFingerprint>> {
232        indexing::file_index::fingerprints_for_scope(self, source_scope)
233    }
234
235    fn code_file_candidate_paths_for_scope(
236        &self,
237        source_scope: String,
238        path_filters: Vec<String>,
239        language_filters: Vec<String>,
240        exclude_generated: bool,
241        limit: usize,
242    ) -> StorageFuture<'_, Vec<String>> {
243        indexing::file_index::candidate_paths_for_scope(
244            self,
245            source_scope,
246            path_filters,
247            language_filters,
248            exclude_generated,
249            limit,
250        )
251    }
252
253    fn code_file_candidate_paths_for_query_scope(
254        &self,
255        source_scope: String,
256        query: String,
257        path_filters: Vec<String>,
258        language_filters: Vec<String>,
259        exclude_generated: bool,
260        limit: usize,
261    ) -> StorageFuture<'_, Vec<String>> {
262        indexing::file_index::candidate_paths_for_query_scope(
263            self,
264            source_scope,
265            query,
266            path_filters,
267            language_filters,
268            exclude_generated,
269            limit,
270        )
271    }
272
273    fn apply_code_index_snapshot(
274        &self,
275        snapshot: CodeIndexSnapshot,
276    ) -> StorageFuture<'_, CodeIndexSummary> {
277        indexing::lifecycle::apply_snapshot(self, snapshot)
278    }
279
280    fn clear_code_workspace_state(
281        &self,
282        repository_id: String,
283        source_scope: String,
284    ) -> StorageFuture<'_, ()> {
285        indexing::lifecycle::clear_workspace(self, repository_id, source_scope)
286    }
287    fn begin_code_index_session(
288        &self,
289        session: CodeIndexSession,
290    ) -> StorageFuture<'_, CodeIndexCheckpoint> {
291        indexing::lifecycle::begin_session(self, session)
292    }
293
294    fn apply_code_index_batch(
295        &self,
296        batch: CodeIndexBatch,
297    ) -> StorageFuture<'_, CodeIndexCheckpoint> {
298        indexing::lifecycle::apply_batch(self, batch)
299    }
300
301    fn finalize_code_index_session(
302        &self,
303        session: CodeIndexSession,
304    ) -> StorageFuture<'_, CodeIndexSummary> {
305        indexing::lifecycle::finalize_session(self, session)
306    }
307
308    fn search_code(
309        &self,
310        request: CodeRetrievalRequest,
311    ) -> StorageFuture<'_, Vec<CodeRetrievalHit>> {
312        let this = self.clone();
313        Box::pin(async move {
314            if let Some(shard) = repository_store_for_selector(
315                &this.control,
316                &this.catalog,
317                request.repository.repository.clone(),
318            )
319            .await?
320            {
321                return match shard.search_code(request.clone()).await {
322                    Ok(hits) => Ok(hits),
323                    Err(error) if is_missing_code_scope_error(&error) => {
324                        this.control.search_code(request).await
325                    }
326                    Err(error) => Err(error),
327                };
328            }
329            this.control.search_code(request).await
330        })
331    }
332
333    fn search_code_feature_flags(
334        &self,
335        request: CodeFeatureFlagRequest,
336    ) -> StorageFuture<'_, Vec<CodeFeatureFlagGraph>> {
337        let this = self.clone();
338        Box::pin(async move {
339            if let Some(shard) = repository_store_for_selector(
340                &this.control,
341                &this.catalog,
342                request.repository.repository.clone(),
343            )
344            .await?
345            {
346                return match shard.search_code_feature_flags(request.clone()).await {
347                    Ok(flags) => Ok(flags),
348                    Err(error) if is_missing_code_scope_error(&error) => {
349                        this.control.search_code_feature_flags(request).await
350                    }
351                    Err(error) => Err(error),
352                };
353            }
354            this.control.search_code_feature_flags(request).await
355        })
356    }
357
358    fn search_code_feature_flags_scope(
359        &self,
360        source_scope: String,
361        request: CodeFeatureFlagRequest,
362    ) -> StorageFuture<'_, Vec<CodeFeatureFlagGraph>> {
363        routing::search_code_feature_flags_scope(self.clone(), source_scope, request)
364    }
365
366    fn search_code_scope(
367        &self,
368        source_scope: String,
369        request: CodeRetrievalRequest,
370    ) -> StorageFuture<'_, Vec<CodeRetrievalHit>> {
371        routing::search_code_scope(self.clone(), source_scope, request)
372    }
373
374    fn analyze_code_impact(
375        &self,
376        request: crate::domain::CodeImpactRequest,
377        changes: CodeImpactChanges,
378    ) -> StorageFuture<'_, Vec<CodeRetrievalHit>> {
379        let this = self.clone();
380        Box::pin(async move {
381            if let Some(shard) = repository_store_for_selector(
382                &this.control,
383                &this.catalog,
384                request.repository.repository.clone(),
385            )
386            .await?
387            {
388                return match shard
389                    .analyze_code_impact(request.clone(), changes.clone())
390                    .await
391                {
392                    Ok(hits) => Ok(hits),
393                    Err(error) if is_missing_code_scope_error(&error) => {
394                        this.control.analyze_code_impact(request, changes).await
395                    }
396                    Err(error) => Err(error),
397                };
398            }
399            this.control.analyze_code_impact(request, changes).await
400        })
401    }
402
403    fn analyze_code_impact_scope(
404        &self,
405        source_scope: String,
406        request: crate::domain::CodeImpactRequest,
407        changes: CodeImpactChanges,
408    ) -> StorageFuture<'_, Vec<CodeRetrievalHit>> {
409        routing::analyze_code_impact_scope(self.clone(), source_scope, request, changes)
410    }
411
412    fn codebase_view_snapshot(
413        &self,
414        source_scope: String,
415        request: crate::domain::CodebaseViewRequest,
416        row_limit: usize,
417    ) -> StorageFuture<'_, crate::domain::CodebaseViewSnapshot> {
418        routing::codebase_view_snapshot(self.clone(), source_scope, request, row_limit)
419    }
420
421    fn code_repository_totals(&self) -> StorageFuture<'_, CodeRepositoryTotals> {
422        let this = self.clone();
423        Box::pin(async move { totals::code_repository_totals(this.control, this.catalog).await })
424    }
425
426    fn code_repository_report(
427        &self,
428        repository: String,
429    ) -> StorageFuture<'_, CodeRepositoryReport> {
430        let this = self.clone();
431        Box::pin(async move {
432            if let Some(shard) =
433                repository_store_for_selector(&this.control, &this.catalog, repository.clone())
434                    .await?
435            {
436                return shard.code_repository_report(repository).await;
437            }
438            this.control.code_repository_report(repository).await
439        })
440    }
441
442    fn code_repository_scope_symbol_generation_counts(
443        &self,
444        source_scope: String,
445    ) -> StorageFuture<'_, CodeSymbolGenerationCounts> {
446        let this = self.clone();
447        Box::pin(async move {
448            totals::scope_symbol_generation_counts(this.control, this.catalog, source_scope).await
449        })
450    }
451
452    fn refresh_software_global_projection(
453        &self,
454        source_scope: String,
455    ) -> StorageFuture<'_, SoftwareGlobalProjection> {
456        let this = self.clone();
457        Box::pin(async move {
458            if let Some(shard) = source_scope_store(&this.catalog, source_scope.clone()).await? {
459                return shard.refresh_software_global_projection(source_scope).await;
460            }
461            this.control
462                .refresh_software_global_projection(source_scope)
463                .await
464        })
465    }
466
467    fn software_global_projection(
468        &self,
469        request: SoftwareGlobalRequest,
470    ) -> StorageFuture<'_, SoftwareGlobalProjection> {
471        let this = self.clone();
472        Box::pin(async move {
473            if let Some(shard) = repository_store_for_selector(
474                &this.control,
475                &this.catalog,
476                request.repository.repository.clone(),
477            )
478            .await?
479            {
480                return match shard.software_global_projection(request.clone()).await {
481                    Ok(projection) => Ok(projection),
482                    Err(error) if is_missing_code_scope_error(&error) => {
483                        this.control.software_global_projection(request).await
484                    }
485                    Err(error) => Err(error),
486                };
487            }
488            this.control.software_global_projection(request).await
489        })
490    }
491
492    fn software_global_projection_for_scope(
493        &self,
494        source_scope: String,
495        request: SoftwareGlobalRequest,
496    ) -> StorageFuture<'_, SoftwareGlobalProjection> {
497        let this = self.clone();
498        Box::pin(async move {
499            if let Some(shard) = source_scope_store(&this.catalog, source_scope.clone()).await? {
500                return shard
501                    .software_global_projection_for_scope(source_scope, request)
502                    .await;
503            }
504            this.control
505                .software_global_projection_for_scope(source_scope, request)
506                .await
507        })
508    }
509
510    fn create_code_repository_set(
511        &self,
512        seed: CodeRepositorySetSeed,
513    ) -> StorageFuture<'_, CodeRepositorySet> {
514        self.control.create_code_repository_set(seed)
515    }
516
517    fn add_code_repository_set_member(
518        &self,
519        seed: CodeRepositorySetMemberSeed,
520    ) -> StorageFuture<'_, CodeRepositorySetMember> {
521        self.control.add_code_repository_set_member(seed)
522    }
523
524    fn remove_code_repository_set_member(
525        &self,
526        set_alias: String,
527        repository_alias: String,
528    ) -> StorageFuture<'_, CodeRepositorySetMember> {
529        self.control
530            .remove_code_repository_set_member(set_alias, repository_alias)
531    }
532
533    fn code_repository_set(
534        &self,
535        set_alias: String,
536    ) -> StorageFuture<'_, Option<CodeRepositorySet>> {
537        self.control.code_repository_set(set_alias)
538    }
539
540    fn code_repository_set_status(
541        &self,
542        set_alias: String,
543    ) -> StorageFuture<'_, Option<CodeRepositorySetStatus>> {
544        self.control.code_repository_set_status(set_alias)
545    }
546
547    fn refresh_code_repository_set_overlay(
548        &self,
549        set_alias: String,
550        _now_ms: u64,
551    ) -> StorageFuture<'_, CodeRepositorySetRefreshSummary> {
552        Box::pin(async move {
553            Err(StorageError::InvalidInput(format!(
554                "repository set overlay refresh for '{set_alias}' requires the single_sqlite topology until cross-shard import/export aggregation is implemented"
555            )))
556        })
557    }
558
559    fn code_repository_set_cross_edges(
560        &self,
561        set_id: String,
562    ) -> StorageFuture<'_, Vec<CodeRepositoryCrossEdge>> {
563        self.control.code_repository_set_cross_edges(set_id)
564    }
565
566    fn code_repository_set_cross_edges_for_selector(
567        &self,
568        set_id: String,
569        selector: crate::storage::CodeRepositorySetEdgeSelector,
570    ) -> StorageFuture<'_, Vec<CodeRepositoryCrossEdge>> {
571        self.control
572            .code_repository_set_cross_edges_for_selector(set_id, selector)
573    }
574
575    fn queue_code_repository_set_refresh_task(
576        &self,
577        task: CodeRepositorySetRefreshTaskSeed,
578    ) -> StorageFuture<'_, crate::domain::CodeRepositorySetRefreshTaskRecord> {
579        Box::pin(async move {
580            Err(StorageError::InvalidInput(format!(
581                "repository set overlay refresh task for '{}' requires the single_sqlite topology until cross-shard import/export aggregation is implemented",
582                task.set_alias
583            )))
584        })
585    }
586
587    fn claim_code_repository_set_refresh_task(
588        &self,
589        request: CodeRepositorySetRefreshTaskClaimRequest,
590    ) -> StorageFuture<'_, Option<crate::domain::CodeRepositorySetRefreshTaskRecord>> {
591        self.control.claim_code_repository_set_refresh_task(request)
592    }
593
594    fn complete_code_repository_set_refresh_task(
595        &self,
596        request: CodeRepositorySetRefreshTaskCompletion,
597    ) -> StorageFuture<'_, crate::domain::CodeRepositorySetRefreshTaskRecord> {
598        self.control
599            .complete_code_repository_set_refresh_task(request)
600    }
601
602    fn fail_code_repository_set_refresh_task(
603        &self,
604        request: CodeRepositorySetRefreshTaskFailure,
605    ) -> StorageFuture<'_, crate::domain::CodeRepositorySetRefreshTaskRecord> {
606        self.control.fail_code_repository_set_refresh_task(request)
607    }
608}
609
610#[cfg(test)]
611#[path = "mod_tests.rs"]
612mod tests;