Skip to main content

relay_knowledge/application/code_repository/indexing/
worker.rs

1//! One-shot durable index-worker claim, completion, retry, and recovery.
2
3use 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    /// Runs one queued code index task under a lease.
29    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    /// Recovers code-index worker leases that belonged to exited service processes.
181    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}