Skip to main content

relay_knowledge/application/code_repository/indexing/
mod.rs

1mod fast_path;
2mod queue;
3mod state;
4mod task;
5mod tasks;
6
7use crate::{
8    api::{
9        ApiError, ApiMetadata, CodeRepositoryIndexResponse, CodeRepositoryIndexStartResponse,
10        CodeRepositoryScopePreviewResponse, RequestContext,
11    },
12    code::{
13        build_index_snapshot_with_workspace_detection,
14        prepare_full_index_plan_with_workspace_detection, preview_repository_scope,
15    },
16    domain::{
17        CodeIndexMode, CodeIndexRequest, CodeIndexResourceBudget, CodeRepositorySelector,
18        CodeRepositoryStatus,
19    },
20};
21
22use crate::application::service::RelayKnowledgeService;
23
24use self::{
25    fast_path::fresh_full_index_response,
26    queue::queue_worktree_overlay_index_task,
27    state::{
28        RETAIN_RECENT_CODE_SCOPES, active_full_index_task_for_request, index_start_from_completed,
29        previous_index_state_for_index, requested_index_ref_for_response,
30    },
31    task::{
32        CODE_INDEX_TASK_LEASE_MS, CODE_INDEX_TASK_MAX_ATTEMPTS, CODE_INDEX_TASK_RETRY_BACKOFF_MS,
33        CodeIndexTaskLeaseContext, code_index_worker_lease_owner,
34        recover_orphaned_code_index_task_leases, refresh_code_index_task_lease,
35    },
36};
37use super::{
38    blocking::run_blocking_code,
39    clock::now_millis,
40    errors::storage_api_error,
41    repository::{registration_from_status, required_code_repository},
42    worktree_ref::pending_worktree_overlay_base_commit,
43};
44
45pub(super) use task::recover_code_index_task_leases;
46
47const PARSED_BATCH_QUEUE_CAPACITY: usize = 2;
48
49impl RelayKnowledgeService {
50    /// Builds or updates the tree-sitter code index for a registered repository.
51    pub async fn index_code_repository(
52        &self,
53        request: CodeIndexRequest,
54        context: RequestContext,
55    ) -> Result<CodeRepositoryIndexResponse, ApiError> {
56        self.index_code_repository_inner(request, context, None)
57            .await
58    }
59
60    async fn index_code_repository_inner(
61        &self,
62        request: CodeIndexRequest,
63        context: RequestContext,
64        task_lease: Option<CodeIndexTaskLeaseContext>,
65    ) -> Result<CodeRepositoryIndexResponse, ApiError> {
66        let store = self.store().await.map_err(storage_api_error)?;
67        let status = required_code_repository(&store, &request.repository.repository).await?;
68        if let Some(response) =
69            fresh_full_index_response(&store, &status, &request, &context).await?
70        {
71            return Ok(response);
72        }
73        let requested_ref = requested_index_ref_for_response(&request);
74        let registration = registration_from_status(&status);
75        let selector = request.repository.clone();
76        let summary = if request.mode == CodeIndexMode::Full {
77            self.apply_full_code_index(
78                &store,
79                registration,
80                selector,
81                request.workspace_detection.clone(),
82                CodeIndexResourceBudget::default(),
83                task_lease,
84            )
85            .await?
86        } else {
87            let previous = previous_index_state_for_index(&store, &status, &request).await?;
88            let mode = request.mode;
89            let workspace_detection = request.workspace_detection.clone();
90            let snapshot = run_blocking_code(move || {
91                build_index_snapshot_with_workspace_detection(
92                    &registration,
93                    &selector,
94                    mode,
95                    previous.fingerprints,
96                    previous.base_resolved_commit_sha,
97                    &workspace_detection,
98                )
99            })
100            .await?;
101            store
102                .apply_code_index_snapshot(snapshot)
103                .await
104                .map_err(storage_api_error)?
105        };
106        let status = store
107            .code_repository_status(summary.repository_id.clone())
108            .await
109            .map_err(storage_api_error)?
110            .ok_or_else(|| ApiError::storage_unavailable("code repository status is missing"))?;
111        let software_projection = store
112            .refresh_software_global_projection(summary.source_scope.clone())
113            .await
114            .map_err(storage_api_error)?;
115        let graph_version = store
116            .current_graph_version()
117            .await
118            .map_err(storage_api_error)?;
119        let degraded_reason = status
120            .degraded_reason
121            .clone()
122            .or(software_projection.status.last_error.clone());
123        let status = CodeRepositoryStatus {
124            degraded_reason,
125            ..status
126        };
127        let _ = self.refresh_watched_code_repository(&status).await;
128
129        Ok(CodeRepositoryIndexResponse {
130            metadata: ApiMetadata::graph_only(&context, graph_version),
131            scope: crate::api::CodeRepositoryScopeMetadata::from_status(
132                &status,
133                &request.repository,
134                requested_ref,
135            ),
136            summary,
137            status,
138        })
139    }
140
141    async fn apply_full_code_index(
142        &self,
143        store: &std::sync::Arc<dyn crate::storage::KnowledgeStore>,
144        registration: crate::domain::CodeRepositoryRegistration,
145        selector: CodeRepositorySelector,
146        workspace_detection: crate::domain::CodeWorkspaceDetectionConfig,
147        resource_budget: CodeIndexResourceBudget,
148        task_lease: Option<CodeIndexTaskLeaseContext>,
149    ) -> Result<crate::domain::CodeIndexSummary, ApiError> {
150        let plan = run_blocking_code(move || {
151            prepare_full_index_plan_with_workspace_detection(
152                registration,
153                selector,
154                resource_budget,
155                &workspace_detection,
156            )
157        })
158        .await?;
159        let session = plan.session();
160        store
161            .begin_code_index_session(session.clone())
162            .await
163            .map_err(storage_api_error)?;
164        refresh_code_index_task_lease(store, task_lease.as_ref()).await?;
165        let (batch_sender, mut batch_receiver) =
166            tokio::sync::mpsc::channel(PARSED_BATCH_QUEUE_CAPACITY);
167        let parser = tokio::spawn(run_blocking_code(move || {
168            let mut plan = plan;
169            loop {
170                let (next_plan, batch) = plan.parse_next_batch()?;
171                plan = next_plan;
172                let Some(batch) = batch else {
173                    return Ok(());
174                };
175                if batch_sender.blocking_send(batch).is_err() {
176                    return Ok(());
177                }
178            }
179        }));
180        let writer_result = async {
181            while let Some(batch) = batch_receiver.recv().await {
182                store
183                    .apply_code_index_batch(batch)
184                    .await
185                    .map_err(storage_api_error)?;
186                refresh_code_index_task_lease(store, task_lease.as_ref()).await?;
187            }
188            Ok::<(), ApiError>(())
189        }
190        .await;
191        drop(batch_receiver);
192        let parser_result = parser
193            .await
194            .map_err(|error| ApiError::storage_unavailable(error.to_string()))?;
195        writer_result?;
196        parser_result?;
197
198        refresh_code_index_task_lease(store, task_lease.as_ref()).await?;
199        let summary = store
200            .finalize_code_index_session(session)
201            .await
202            .map_err(storage_api_error)?;
203        refresh_code_index_task_lease(store, task_lease.as_ref()).await?;
204
205        Ok(summary)
206    }
207
208    /// Starts a repository index request, queueing cold full indexes for background execution.
209    pub async fn start_code_repository_index(
210        &self,
211        request: CodeIndexRequest,
212        context: RequestContext,
213    ) -> Result<CodeRepositoryIndexStartResponse, ApiError> {
214        let store = self.store().await.map_err(storage_api_error)?;
215        let status = required_code_repository(&store, &request.repository.repository).await?;
216        if let Some(response) =
217            fresh_full_index_response(&store, &status, &request, &context).await?
218        {
219            return Ok(index_start_from_completed(response, None));
220        }
221        if matches!(request.mode, CodeIndexMode::Incremental { .. }) {
222            let response = self.index_code_repository(request, context).await?;
223            return Ok(index_start_from_completed(response, None));
224        }
225        recover_code_index_task_leases(&store, now_millis()).await?;
226        if request.mode == CodeIndexMode::WorktreeOverlay {
227            let requested_ref = requested_index_ref_for_response(&request);
228            let task = queue_worktree_overlay_index_task(&store, &status, &request).await?;
229            return self
230                .index_start_response_from_task(&store, status, task, requested_ref, &context)
231                .await;
232        }
233        let payload_json = serde_json::to_string(&request)
234            .map_err(|error| ApiError::invalid_argument(error.to_string()))?;
235        if let Some(active_task) =
236            active_full_index_task_for_request(&store, &status, &request, &payload_json).await?
237        {
238            return self
239                .index_start_response_from_task(
240                    &store,
241                    status,
242                    active_task,
243                    request.repository.ref_selector,
244                    &context,
245                )
246                .await;
247        }
248
249        let registration = registration_from_status(&status);
250        let selector = request.repository.clone();
251        let workspace_detection = request.workspace_detection.clone();
252        let workspace_detection_json = serde_json::to_string(&workspace_detection)
253            .map_err(|error| ApiError::invalid_argument(error.to_string()))?;
254        let resource_budget = CodeIndexResourceBudget::default();
255        let plan = run_blocking_code(move || {
256            prepare_full_index_plan_with_workspace_detection(
257                registration,
258                selector,
259                resource_budget,
260                &workspace_detection,
261            )
262        })
263        .await?;
264        let session = plan.session();
265        let input_fingerprint = format!(
266            "full:{}:{}:{}:{}",
267            session.repository_id,
268            session.tree_hash,
269            session.source_scope,
270            workspace_detection_json
271        );
272        if let Some(active_task) = store
273            .active_code_index_task(session.repository_id.clone())
274            .await
275            .map_err(storage_api_error)?
276            && active_task.state.is_unfinished()
277            && active_task.input_fingerprint == input_fingerprint
278        {
279            return self
280                .index_start_response_from_task(
281                    &store,
282                    status,
283                    active_task,
284                    request.repository.ref_selector,
285                    &context,
286                )
287                .await;
288        }
289        let task = store
290            .queue_code_index_task(crate::storage::CodeIndexTaskSeed {
291                repository_id: session.repository_id.clone(),
292                alias: status.alias.clone(),
293                ref_selector: request.repository.ref_selector.clone(),
294                resolved_commit_sha: session.resolved_commit_sha.clone(),
295                tree_hash: session.tree_hash.clone(),
296                source_scope: session.source_scope.clone(),
297                path_filters: session.path_filters.clone(),
298                language_filters: session.language_filters.clone(),
299                mode: request.mode.clone(),
300                input_fingerprint,
301                resource_budget: session.resource_budget,
302                payload_json,
303                now_ms: now_millis(),
304            })
305            .await
306            .map_err(storage_api_error)?;
307        self.index_start_response_from_task(
308            &store,
309            status,
310            task,
311            request.repository.ref_selector,
312            &context,
313        )
314        .await
315    }
316
317    async fn index_start_response_from_task(
318        &self,
319        store: &std::sync::Arc<dyn crate::storage::KnowledgeStore>,
320        fallback_status: CodeRepositoryStatus,
321        task: crate::domain::CodeIndexTaskRecord,
322        requested_ref: String,
323        context: &RequestContext,
324    ) -> Result<CodeRepositoryIndexStartResponse, ApiError> {
325        let checkpoint = store
326            .code_index_checkpoint(task.source_scope.clone())
327            .await
328            .map_err(storage_api_error)?;
329        let graph_version = store
330            .current_graph_version()
331            .await
332            .map_err(storage_api_error)?;
333        let status = store
334            .code_repository_status(task.repository_id.clone())
335            .await
336            .map_err(storage_api_error)?
337            .unwrap_or(fallback_status);
338
339        Ok(CodeRepositoryIndexStartResponse {
340            metadata: ApiMetadata::graph_only(context, graph_version),
341            scope: crate::api::CodeRepositoryScopeMetadata::from_index_task(&task, requested_ref),
342            summary: None,
343            status,
344            task: Some(task),
345            checkpoint,
346        })
347    }
348
349    /// Runs one queued code index task under a lease.
350    pub async fn run_code_index_task_once(
351        &self,
352        task_id: Option<String>,
353        context: RequestContext,
354    ) -> Result<Option<crate::domain::CodeIndexTaskRecord>, ApiError> {
355        let store = self.store().await.map_err(storage_api_error)?;
356        let lease_owner = code_index_worker_lease_owner();
357        let Some(task) = store
358            .claim_code_index_task(crate::storage::CodeIndexTaskClaimRequest {
359                task_id,
360                lease_owner: lease_owner.clone(),
361                lease_duration_ms: CODE_INDEX_TASK_LEASE_MS,
362                max_attempts: CODE_INDEX_TASK_MAX_ATTEMPTS,
363                now_ms: now_millis(),
364            })
365            .await
366            .map_err(storage_api_error)?
367        else {
368            return Ok(None);
369        };
370        let mut request = match serde_json::from_str::<CodeIndexRequest>(&task.payload_json) {
371            Ok(request) => request,
372            Err(error) => {
373                let message = format!(
374                    "code index task '{}' payload is invalid: {error}",
375                    task.task_id
376                );
377                let _ = store
378                    .fail_code_index_task(crate::storage::CodeIndexTaskFailure {
379                        task_id: task.task_id,
380                        lease_owner,
381                        attempt_count: task.attempt_count,
382                        error_kind: "task_payload".to_owned(),
383                        error_message: message.clone(),
384                        retry_backoff_ms: CODE_INDEX_TASK_RETRY_BACKOFF_MS,
385                        max_attempts: CODE_INDEX_TASK_MAX_ATTEMPTS,
386                        now_ms: now_millis(),
387                    })
388                    .await;
389                return Err(ApiError::invalid_argument(message));
390            }
391        };
392        if task.mode == CodeIndexMode::WorktreeOverlay {
393            if let Some(base_commit) =
394                pending_worktree_overlay_base_commit(&task.resolved_commit_sha)
395            {
396                request.repository.ref_selector = base_commit.to_owned();
397            }
398        } else {
399            request.repository.ref_selector = task.resolved_commit_sha.clone();
400        }
401        let lease_context = CodeIndexTaskLeaseContext {
402            task_id: task.task_id.clone(),
403            lease_owner: lease_owner.clone(),
404            attempt_count: task.attempt_count,
405            lease_duration_ms: CODE_INDEX_TASK_LEASE_MS,
406        };
407        let result = self
408            .index_code_repository_inner(request, context, Some(lease_context.clone()))
409            .await;
410        match result {
411            Ok(response) => {
412                refresh_code_index_task_lease(&store, Some(&lease_context)).await?;
413                let completed = store
414                    .complete_code_index_task(crate::storage::CodeIndexTaskCompletion {
415                        task_id: task.task_id.clone(),
416                        lease_owner,
417                        attempt_count: task.attempt_count,
418                        now_ms: now_millis(),
419                    })
420                    .await
421                    .map_err(storage_api_error)?;
422                let _ = store
423                    .prune_code_repository_scopes(crate::storage::CodeScopeRetentionRequest {
424                        repository_id: response.summary.repository_id,
425                        active_scope: response.summary.source_scope,
426                        retain_recent_successful_scopes: RETAIN_RECENT_CODE_SCOPES,
427                    })
428                    .await;
429                Ok(Some(completed))
430            }
431            Err(error) => {
432                let _ = store
433                    .fail_code_index_task(crate::storage::CodeIndexTaskFailure {
434                        task_id: task.task_id,
435                        lease_owner,
436                        attempt_count: task.attempt_count,
437                        error_kind: "code_index".to_owned(),
438                        error_message: error.message.clone(),
439                        retry_backoff_ms: CODE_INDEX_TASK_RETRY_BACKOFF_MS,
440                        max_attempts: CODE_INDEX_TASK_MAX_ATTEMPTS,
441                        now_ms: now_millis(),
442                    })
443                    .await;
444                Err(error)
445            }
446        }
447    }
448
449    /// Recovers code-index worker leases that belonged to exited service processes.
450    pub(crate) async fn recover_orphaned_code_index_tasks_on_startup(
451        &self,
452    ) -> Result<usize, ApiError> {
453        let store = self.store().await.map_err(storage_api_error)?;
454        recover_orphaned_code_index_task_leases(
455            &store,
456            now_millis(),
457            &self.runtime.process.windows_tasklist_command,
458        )
459        .await
460    }
461
462    /// Previews the effective code repository indexing scope without writing rows.
463    pub async fn preview_code_repository_scope(
464        &self,
465        request: CodeIndexRequest,
466        context: RequestContext,
467    ) -> Result<CodeRepositoryScopePreviewResponse, ApiError> {
468        let store = self.store().await.map_err(storage_api_error)?;
469        let status = required_code_repository(&store, &request.repository.repository).await?;
470        let registration = registration_from_status(&status);
471        let selector = request.repository.clone();
472        let preview =
473            run_blocking_code(move || preview_repository_scope(&registration, &selector)).await?;
474        let graph_version = store
475            .current_graph_version()
476            .await
477            .map_err(storage_api_error)?;
478        Ok(CodeRepositoryScopePreviewResponse {
479            metadata: ApiMetadata::graph_only(&context, graph_version),
480            scope: crate::api::CodeRepositoryScopeMetadata::from_status(
481                &status,
482                &request.repository,
483                request.repository.ref_selector.clone(),
484            ),
485            preview,
486        })
487    }
488}