relay_knowledge/application/code_repository/indexing/
start.rs1use 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 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}