relay_knowledge/application/code_repository/indexing/
worker.rs1use crate::{
4 api::{ApiError, CodeRepositoryIndexResponse, RequestContext},
5 application::service::RelayKnowledgeService,
6 domain::{CodeIndexMode, CodeIndexRequest, CodeIndexTaskRecord},
7 storage::{
8 CodeIndexTaskClaimRequest, CodeIndexTaskCompletion, CodeIndexTaskFailure,
9 CodeScopeRetentionRequest,
10 },
11};
12
13use super::{
14 super::{
15 clock::now_millis, errors::storage_api_error,
16 worktree_ref::pending_worktree_overlay_base_commit,
17 },
18 state::RETAIN_RECENT_CODE_SCOPES,
19 task::{
20 CODE_INDEX_TASK_LEASE_MS, CODE_INDEX_TASK_MAX_ATTEMPTS, CODE_INDEX_TASK_RETRY_BACKOFF_MS,
21 CodeIndexTaskLeaseContext, code_index_task_failure_disposition,
22 code_index_worker_lease_owner, recover_orphaned_code_index_task_leases,
23 refresh_code_index_task_lease,
24 },
25};
26
27impl RelayKnowledgeService {
28 pub async fn run_code_index_task_once(
30 &self,
31 task_id: Option<String>,
32 context: RequestContext,
33 ) -> Result<Option<CodeIndexTaskRecord>, ApiError> {
34 self.run_code_index_task_once_with_response(task_id, context)
35 .await
36 .map(|outcome| outcome.map(|(task, _)| task))
37 }
38
39 pub(crate) async fn run_code_index_task_once_with_response(
40 &self,
41 task_id: Option<String>,
42 context: RequestContext,
43 ) -> Result<Option<(CodeIndexTaskRecord, CodeRepositoryIndexResponse)>, ApiError> {
44 let store = self.store().await.map_err(storage_api_error)?;
45 let lease_owner = code_index_worker_lease_owner();
46 let Some(task) = store
47 .claim_code_index_task(CodeIndexTaskClaimRequest {
48 task_id,
49 lease_owner: lease_owner.clone(),
50 lease_duration_ms: CODE_INDEX_TASK_LEASE_MS,
51 max_attempts: CODE_INDEX_TASK_MAX_ATTEMPTS,
52 now_ms: now_millis(),
53 })
54 .await
55 .map_err(storage_api_error)?
56 else {
57 return Ok(None);
58 };
59 let mut request = match serde_json::from_str::<CodeIndexRequest>(&task.payload_json) {
60 Ok(request) => request,
61 Err(error) => {
62 let message = format!(
63 "code index task '{}' payload is invalid: {error}",
64 task.task_id
65 );
66 let _ = store
67 .fail_code_index_task(CodeIndexTaskFailure {
68 task_id: task.task_id,
69 lease_owner,
70 attempt_count: task.attempt_count,
71 publication_generation: task.publication_generation,
72 error_kind: "task_payload".to_owned(),
73 error_message: message.clone(),
74 retry_backoff_ms: CODE_INDEX_TASK_RETRY_BACKOFF_MS,
75 max_attempts: CODE_INDEX_TASK_MAX_ATTEMPTS,
76 now_ms: now_millis(),
77 })
78 .await;
79 return Err(ApiError::invalid_argument(message));
80 }
81 };
82 if task.mode == CodeIndexMode::WorktreeOverlay {
83 if let Some(base_commit) =
84 pending_worktree_overlay_base_commit(&task.resolved_commit_sha)
85 {
86 request.repository.ref_selector = base_commit.to_owned();
87 }
88 } else if task.mode == CodeIndexMode::Full {
89 request.repository.ref_selector = task.resolved_commit_sha.clone();
90 } else {
91 request.repository.ref_selector = task.ref_selector.clone();
92 }
93 let lease_context = CodeIndexTaskLeaseContext {
94 task_id: task.task_id.clone(),
95 lease_owner: lease_owner.clone(),
96 attempt_count: task.attempt_count,
97 lease_duration_ms: CODE_INDEX_TASK_LEASE_MS,
98 publication_fence: crate::domain::CodeIndexPublicationFence {
99 repository_id: task.repository_id.clone(),
100 task_id: task.task_id.clone(),
101 lease_owner: lease_owner.clone(),
102 attempt_count: task.attempt_count,
103 generation: task.publication_generation,
104 },
105 source_scope: task.source_scope.clone(),
106 resolved_commit_sha: task.resolved_commit_sha.clone(),
107 tree_hash: task.tree_hash.clone(),
108 path_filters: task.path_filters.clone(),
109 language_filters: task.language_filters.clone(),
110 resource_budget: task.resource_budget,
111 };
112 let result = self
113 .index_code_repository_inner(request, context, Some(lease_context.clone()))
114 .await;
115 match result {
116 Ok(response) => {
117 refresh_code_index_task_lease(&store, Some(&lease_context)).await?;
118 let completed = store
119 .complete_code_index_task(CodeIndexTaskCompletion {
120 task_id: task.task_id.clone(),
121 lease_owner,
122 attempt_count: task.attempt_count,
123 publication_generation: task.publication_generation,
124 now_ms: now_millis(),
125 })
126 .await
127 .map_err(storage_api_error)?;
128 if let Err(error) = store
129 .run_code_index_post_maintenance(
130 response.summary.repository_id.clone(),
131 response.summary.source_scope.clone(),
132 )
133 .await
134 {
135 tracing::warn!(
136 task_id = %completed.task_id,
137 error = %error,
138 "code index completed but post-index SQLite maintenance did not run"
139 );
140 }
141 if let Err(error) = store
142 .prune_code_repository_scopes(CodeScopeRetentionRequest {
143 repository_id: response.summary.repository_id.clone(),
144 active_scope: response.summary.source_scope.clone(),
145 retain_recent_successful_scopes: RETAIN_RECENT_CODE_SCOPES,
146 repository_retention_cutoff_ms: None,
147 repository_retention_cutoff_generation: None,
148 repository_retention_initial_scope: None,
149 })
150 .await
151 {
152 tracing::warn!(
153 task_id = %completed.task_id,
154 error = %error,
155 "code index published but bounded scope retention did not complete"
156 );
157 }
158 Ok(Some((completed, response)))
159 }
160 Err(error) => {
161 let failure = code_index_task_failure_disposition(&error, task.attempt_count);
162 let _ = store
163 .fail_code_index_task(CodeIndexTaskFailure {
164 task_id: task.task_id,
165 lease_owner,
166 attempt_count: task.attempt_count,
167 publication_generation: task.publication_generation,
168 error_kind: failure.error_kind.to_owned(),
169 error_message: error.message.clone(),
170 retry_backoff_ms: CODE_INDEX_TASK_RETRY_BACKOFF_MS,
171 max_attempts: failure.max_attempts,
172 now_ms: now_millis(),
173 })
174 .await;
175 Err(error)
176 }
177 }
178 }
179
180 pub(crate) async fn recover_orphaned_code_index_tasks_on_startup(
182 &self,
183 ) -> Result<usize, ApiError> {
184 let store = self.store().await.map_err(storage_api_error)?;
185 recover_orphaned_code_index_task_leases(
186 &store,
187 now_millis(),
188 &self.runtime.process.windows_tasklist_command,
189 )
190 .await
191 }
192}