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