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