Skip to main content

relay_knowledge/storage/partitioned/control_plane/
mod.rs

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