Skip to main content

relay_knowledge/application/code_repository/indexing/
start.rs

1//! Durable task admission and historical-reuse planning for index requests.
2
3use crate::{
4    api::{ApiError, CodeRepositoryIndexStartResponse, RequestContext},
5    application::service::RelayKnowledgeService,
6    code::prepare_full_index_plan_with_workspace_detection,
7    domain::{CodeIndexMode, CodeIndexRequest, CodeIndexResourceBudget},
8    storage::CodeIndexTaskSeed,
9};
10
11use super::{
12    super::{
13        blocking::run_blocking_code,
14        clock::now_millis,
15        errors::storage_api_error,
16        repository::{registration_from_status, required_code_repository},
17    },
18    fast_path::fresh_full_index_response,
19    queue::{
20        index_start_response_from_task, queue_incremental_index_task,
21        queue_worktree_overlay_index_task,
22    },
23    state::{
24        FullIndexReusePlan, active_full_index_task_for_request,
25        historical_reuse_base_became_unavailable, index_start_from_completed,
26        plan_full_index_reuse, requested_index_ref_for_response,
27    },
28    task::recover_code_index_task_leases,
29};
30
31impl RelayKnowledgeService {
32    /// Starts a repository index request under the durable single-writer queue.
33    pub async fn start_code_repository_index(
34        &self,
35        request: CodeIndexRequest,
36        context: RequestContext,
37    ) -> Result<CodeRepositoryIndexStartResponse, ApiError> {
38        let store = self.store().await.map_err(storage_api_error)?;
39        let status =
40            required_code_repository(store.as_ref(), &request.repository.repository).await?;
41        if let Some(response) =
42            fresh_full_index_response(&store, &status, &request, &context).await?
43        {
44            return Ok(index_start_from_completed(response, None));
45        }
46        recover_code_index_task_leases(&store, now_millis()).await?;
47        if matches!(&request.mode, CodeIndexMode::Incremental { .. }) {
48            let requested_ref = requested_index_ref_for_response(&request);
49            let task = queue_incremental_index_task(&store, &status, &request).await?;
50            return index_start_response_from_task(&store, status, task, requested_ref, &context)
51                .await;
52        }
53        if request.mode == CodeIndexMode::WorktreeOverlay {
54            let requested_ref = requested_index_ref_for_response(&request);
55            let task = queue_worktree_overlay_index_task(&store, &status, &request).await?;
56            return index_start_response_from_task(&store, status, task, requested_ref, &context)
57                .await;
58        }
59        match if request.reuse_historical {
60            plan_full_index_reuse(&store, &status, &request).await?
61        } else {
62            FullIndexReusePlan::Full
63        } {
64            FullIndexReusePlan::ActiveTask(task) => {
65                return index_start_response_from_task(
66                    &store,
67                    status,
68                    *task,
69                    request.repository.ref_selector,
70                    &context,
71                )
72                .await;
73            }
74            FullIndexReusePlan::Incremental(incremental_request) => {
75                let requested_ref = request.repository.ref_selector.clone();
76                match queue_incremental_index_task(&store, &status, &incremental_request).await {
77                    Ok(task) => {
78                        return index_start_response_from_task(
79                            &store,
80                            status,
81                            task,
82                            requested_ref,
83                            &context,
84                        )
85                        .await;
86                    }
87                    Err(error) if historical_reuse_base_became_unavailable(&error) => {}
88                    Err(error) => return Err(error),
89                }
90            }
91            FullIndexReusePlan::Full => {}
92        }
93        let payload_json = serde_json::to_string(&request)
94            .map_err(|error| ApiError::invalid_argument(error.to_string()))?;
95        if let Some(active_task) =
96            active_full_index_task_for_request(&store, &status, &request, &payload_json).await?
97        {
98            return index_start_response_from_task(
99                &store,
100                status,
101                active_task,
102                request.repository.ref_selector,
103                &context,
104            )
105            .await;
106        }
107
108        let registration = registration_from_status(&status);
109        let selector = request.repository.clone();
110        let workspace_detection = request.workspace_detection.clone();
111        let workspace_detection_json = serde_json::to_string(&workspace_detection)
112            .map_err(|error| ApiError::invalid_argument(error.to_string()))?;
113        let resource_budget = CodeIndexResourceBudget::default();
114        let plan = run_blocking_code(move || {
115            prepare_full_index_plan_with_workspace_detection(
116                registration,
117                selector,
118                resource_budget,
119                &workspace_detection,
120            )
121        })
122        .await?;
123        let session = plan.session();
124        let input_fingerprint = format!(
125            "full:{}:{}:{}:{}",
126            session.repository_id,
127            session.tree_hash,
128            session.source_scope,
129            workspace_detection_json
130        );
131        if let Some(active_task) = store
132            .active_code_index_task(session.repository_id.clone())
133            .await
134            .map_err(storage_api_error)?
135            && active_task.state.is_unfinished()
136            && active_task.input_fingerprint == input_fingerprint
137        {
138            return index_start_response_from_task(
139                &store,
140                status,
141                active_task,
142                request.repository.ref_selector,
143                &context,
144            )
145            .await;
146        }
147        let task = store
148            .queue_code_index_task(CodeIndexTaskSeed {
149                repository_id: session.repository_id.clone(),
150                alias: status.alias.clone(),
151                ref_selector: request.repository.ref_selector.clone(),
152                resolved_commit_sha: session.resolved_commit_sha.clone(),
153                tree_hash: session.tree_hash.clone(),
154                source_scope: session.source_scope.clone(),
155                path_filters: session.path_filters.clone(),
156                language_filters: session.language_filters.clone(),
157                mode: request.mode.clone(),
158                input_fingerprint,
159                resource_budget: session.resource_budget,
160                payload_json,
161                now_ms: now_millis(),
162            })
163            .await
164            .map_err(storage_api_error)?;
165        index_start_response_from_task(
166            &store,
167            status,
168            task,
169            request.repository.ref_selector,
170            &context,
171        )
172        .await
173    }
174}