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 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 ®istration,
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 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 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 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 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(®istration, &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}