Skip to main content

relay_knowledge/storage/partitioned/
control_delegates.rs

1use crate::{
2    domain::{
3        AuditEventRecord, CodeChunkRecord, CodeGraphBatch, CodeGraphCommitReceipt,
4        CodeIndexSnapshot, CodeReferenceRecord, CodeRepositoryStatus, CodeSymbolRecord,
5        CommitReceipt, GraphMutationBatch, GraphVersion, IndexKind, IndexStatus, ProposalState,
6        RetrievalHit, ServiceOperatorStatus, WorkerStatus, WorkerTaskRecord,
7    },
8    storage::{
9        AuditQueryRequest, CodeChunkSearchRequest, CodeGraphStore, CodeReferenceSearchRequest,
10        CodeRepositoryStore, CodeSymbolSearchRequest, FileIndexDiagnostics, FileIndexRoot,
11        FileIndexRootStatus, FileIndexRootUpdate, FileSearchHit, FileSearchRequest,
12        GraphCanvasStorageRequest, GraphCanvasStorageSnapshot, GraphInspection, GraphSearchRequest,
13        GraphStore, HealthStorageSnapshot, IndexCursor, IndexRefreshClaimRequest,
14        IndexRefreshCompletion, IndexRefreshDiagnostics, IndexRefreshFailure,
15        IndexRefreshQueueRequest, IndexRefreshTask, IndexStore, MutationLogEntry, MutationLogStore,
16        NewAuditEvent, NewProposal, ProposalDecision, ProposalListRequest, ServiceOperatorUpdate,
17        StorageError, StorageFuture, WorkerTaskClaimRequest, WorkerTaskCompletion,
18        WorkerTaskFailure, WorkerTaskSeed,
19    },
20};
21
22use super::PartitionedSqliteKnowledgeStore;
23
24pub(super) fn list_code_repositories(
25    store: &PartitionedSqliteKnowledgeStore,
26) -> StorageFuture<'_, Vec<CodeRepositoryStatus>> {
27    let this = store.clone();
28    Box::pin(async move {
29        let control_statuses = this.control.list_code_repositories().await?;
30        let mut statuses = Vec::with_capacity(control_statuses.len());
31        for control_status in control_statuses {
32            let Some(shard) = this
33                .catalog
34                .existing_repository_store(control_status.repository_id.clone())
35                .await?
36            else {
37                statuses.push(control_status);
38                continue;
39            };
40            let Some(mut shard_status) = shard
41                .code_repository_status(control_status.repository_id.clone())
42                .await?
43            else {
44                statuses.push(control_status);
45                continue;
46            };
47            shard_status.alias = control_status.alias;
48            statuses.push(shard_status);
49        }
50        Ok(statuses)
51    })
52}
53
54pub(super) async fn incremental_base_scope(
55    store: &PartitionedSqliteKnowledgeStore,
56    snapshot: &CodeIndexSnapshot,
57) -> Result<Option<String>, StorageError> {
58    if snapshot.full_replace {
59        return Ok(None);
60    }
61    let Some(base_commit) = snapshot.base_resolved_commit_sha.clone() else {
62        return Ok(None);
63    };
64
65    Ok(store
66        .control
67        .code_repository_scope_status(
68            snapshot.repository_id.clone(),
69            base_commit,
70            snapshot.path_filters.clone(),
71            snapshot.language_filters.clone(),
72        )
73        .await?
74        .and_then(|status| status.last_indexed_scope_id))
75}
76
77impl GraphStore for PartitionedSqliteKnowledgeStore {
78    fn commit_mutation_batch(&self, batch: GraphMutationBatch) -> StorageFuture<'_, CommitReceipt> {
79        self.control.commit_mutation_batch(batch)
80    }
81
82    fn inspect_graph(&self) -> StorageFuture<'_, GraphInspection> {
83        let this = self.clone();
84        Box::pin(async move { super::diagnostics::inspect_graph(&this).await })
85    }
86
87    fn health_snapshot(&self, now_ms: u64) -> StorageFuture<'_, HealthStorageSnapshot> {
88        let this = self.clone();
89        Box::pin(async move {
90            let mut snapshot = super::diagnostics::health_snapshot(&this, now_ms).await?;
91            snapshot.repository_code_totals = this.code_repository_totals().await?;
92            Ok(snapshot)
93        })
94    }
95
96    fn graph_canvas(
97        &self,
98        request: GraphCanvasStorageRequest,
99    ) -> StorageFuture<'_, GraphCanvasStorageSnapshot> {
100        self.control.graph_canvas(request)
101    }
102
103    fn search(&self, request: GraphSearchRequest) -> StorageFuture<'_, Vec<RetrievalHit>> {
104        self.control.search(request)
105    }
106
107    fn current_graph_version(&self) -> StorageFuture<'_, GraphVersion> {
108        self.control.current_graph_version()
109    }
110}
111
112impl MutationLogStore for PartitionedSqliteKnowledgeStore {
113    fn read_after(
114        &self,
115        graph_version: GraphVersion,
116        limit: usize,
117    ) -> StorageFuture<'_, Vec<MutationLogEntry>> {
118        self.control.read_after(graph_version, limit)
119    }
120}
121
122impl IndexStore for PartitionedSqliteKnowledgeStore {
123    fn index_statuses(&self) -> StorageFuture<'_, Vec<IndexStatus>> {
124        self.control.index_statuses()
125    }
126
127    fn mark_refresh_complete(
128        &self,
129        kind: IndexKind,
130        graph_version: GraphVersion,
131    ) -> StorageFuture<'_, IndexStatus> {
132        self.control.mark_refresh_complete(kind, graph_version)
133    }
134
135    fn index_cursors(&self) -> StorageFuture<'_, Vec<IndexCursor>> {
136        self.control.index_cursors()
137    }
138
139    fn queue_index_refreshes(
140        &self,
141        request: IndexRefreshQueueRequest,
142    ) -> StorageFuture<'_, IndexRefreshDiagnostics> {
143        self.control.queue_index_refreshes(request)
144    }
145
146    fn claim_index_refresh_task(
147        &self,
148        request: IndexRefreshClaimRequest,
149    ) -> StorageFuture<'_, Option<IndexRefreshTask>> {
150        self.control.claim_index_refresh_task(request)
151    }
152
153    fn complete_index_refresh_task(
154        &self,
155        request: IndexRefreshCompletion,
156    ) -> StorageFuture<'_, IndexRefreshTask> {
157        self.control.complete_index_refresh_task(request)
158    }
159
160    fn fail_index_refresh_task(
161        &self,
162        request: IndexRefreshFailure,
163    ) -> StorageFuture<'_, IndexRefreshTask> {
164        self.control.fail_index_refresh_task(request)
165    }
166
167    fn index_refresh_diagnostics(&self, now_ms: u64) -> StorageFuture<'_, IndexRefreshDiagnostics> {
168        self.control.index_refresh_diagnostics(now_ms)
169    }
170
171    fn queue_worker_tasks(
172        &self,
173        tasks: Vec<WorkerTaskSeed>,
174    ) -> StorageFuture<'_, Vec<WorkerTaskRecord>> {
175        self.control.queue_worker_tasks(tasks)
176    }
177
178    fn worker_statuses(&self) -> StorageFuture<'_, Vec<WorkerStatus>> {
179        self.control.worker_statuses()
180    }
181
182    fn claim_worker_task(
183        &self,
184        request: WorkerTaskClaimRequest,
185    ) -> StorageFuture<'_, Option<WorkerTaskRecord>> {
186        self.control.claim_worker_task(request)
187    }
188
189    fn complete_worker_task(
190        &self,
191        request: WorkerTaskCompletion,
192    ) -> StorageFuture<'_, WorkerTaskRecord> {
193        self.control.complete_worker_task(request)
194    }
195
196    fn fail_worker_task(&self, request: WorkerTaskFailure) -> StorageFuture<'_, WorkerTaskRecord> {
197        self.control.fail_worker_task(request)
198    }
199
200    fn insert_proposal(
201        &self,
202        proposal: NewProposal,
203    ) -> StorageFuture<'_, crate::domain::ProposalRecord> {
204        self.control.insert_proposal(proposal)
205    }
206
207    fn list_proposals(
208        &self,
209        request: ProposalListRequest,
210    ) -> StorageFuture<'_, Vec<crate::domain::ProposalRecord>> {
211        self.control.list_proposals(request)
212    }
213
214    fn proposal_count(&self, state: Option<ProposalState>) -> StorageFuture<'_, usize> {
215        self.control.proposal_count(state)
216    }
217
218    fn proposal_by_id(
219        &self,
220        proposal_id: String,
221    ) -> StorageFuture<'_, Option<crate::domain::ProposalRecord>> {
222        self.control.proposal_by_id(proposal_id)
223    }
224
225    fn proposal_conflicts(
226        &self,
227        proposal_id: String,
228    ) -> StorageFuture<'_, Vec<crate::domain::ProposalConflictRecord>> {
229        self.control.proposal_conflicts(proposal_id)
230    }
231
232    fn decide_proposal(
233        &self,
234        request: ProposalDecision,
235    ) -> StorageFuture<'_, crate::domain::ProposalRecord> {
236        self.control.decide_proposal(request)
237    }
238
239    fn insert_audit_event(&self, event: NewAuditEvent) -> StorageFuture<'_, AuditEventRecord> {
240        self.control.insert_audit_event(event)
241    }
242
243    fn query_audit_events(
244        &self,
245        request: AuditQueryRequest,
246    ) -> StorageFuture<'_, Vec<AuditEventRecord>> {
247        self.control.query_audit_events(request)
248    }
249
250    fn audit_event_count(&self) -> StorageFuture<'_, usize> {
251        self.control.audit_event_count()
252    }
253
254    fn service_operator_status(&self) -> StorageFuture<'_, ServiceOperatorStatus> {
255        self.control.service_operator_status()
256    }
257
258    fn update_service_operator(
259        &self,
260        request: ServiceOperatorUpdate,
261    ) -> StorageFuture<'_, ServiceOperatorStatus> {
262        self.control.update_service_operator(request)
263    }
264
265    fn replace_file_index_root(
266        &self,
267        update: FileIndexRootUpdate,
268    ) -> StorageFuture<'_, FileIndexRootStatus> {
269        self.control.replace_file_index_root(update)
270    }
271
272    fn mark_file_index_roots_unconfigured(
273        &self,
274        active_roots: Vec<FileIndexRoot>,
275        now_ms: u64,
276    ) -> StorageFuture<'_, FileIndexDiagnostics> {
277        self.control
278            .mark_file_index_roots_unconfigured(active_roots, now_ms)
279    }
280
281    fn search_files(&self, request: FileSearchRequest) -> StorageFuture<'_, Vec<FileSearchHit>> {
282        self.control.search_files(request)
283    }
284
285    fn file_index_diagnostics(&self) -> StorageFuture<'_, FileIndexDiagnostics> {
286        self.control.file_index_diagnostics()
287    }
288}
289
290impl CodeGraphStore for PartitionedSqliteKnowledgeStore {
291    fn commit_code_graph_batch(
292        &self,
293        batch: CodeGraphBatch,
294    ) -> StorageFuture<'_, CodeGraphCommitReceipt> {
295        self.control.commit_code_graph_batch(batch)
296    }
297
298    fn search_code_symbols(
299        &self,
300        request: CodeSymbolSearchRequest,
301    ) -> StorageFuture<'_, Vec<CodeSymbolRecord>> {
302        self.control.search_code_symbols(request)
303    }
304
305    fn search_code_references(
306        &self,
307        request: CodeReferenceSearchRequest,
308    ) -> StorageFuture<'_, Vec<CodeReferenceRecord>> {
309        self.control.search_code_references(request)
310    }
311
312    fn search_code_chunks(
313        &self,
314        request: CodeChunkSearchRequest,
315    ) -> StorageFuture<'_, Vec<CodeChunkRecord>> {
316        self.control.search_code_chunks(request)
317    }
318}