Skip to main content

relay_knowledge/application/code_repository/indexing/
tasks.rs

1use std::sync::atomic::Ordering;
2
3use crate::{
4    api::{
5        ApiError, ApiMetadata, CodeIndexWorkerStatus, CodeRepositoryIndexResetResponse,
6        RequestContext,
7    },
8    application::service::RelayKnowledgeService,
9};
10
11use super::super::{
12    clock::now_millis,
13    errors::storage_api_error,
14    repository::{code_status_checkpoint, required_code_repository},
15};
16use super::state::RETAIN_RECENT_CODE_SCOPES;
17use super::task::recover_orphaned_code_index_task_leases;
18
19impl RelayKnowledgeService {
20    /// Runs one bounded, restart-safe retention pass for persistent leftovers.
21    pub(crate) async fn run_code_scope_retention_once(&self) -> Result<bool, ApiError> {
22        let store = self.store().await.map_err(storage_api_error)?;
23        store
24            .schedule_code_repository_retention(
25                self.runtime.workers.code_index_max_indexed_repositories,
26                now_millis(),
27            )
28            .await
29            .map_err(storage_api_error)?;
30        let repository_scan_pending = store
31            .code_repository_retention_scan_pending()
32            .await
33            .map_err(storage_api_error)?;
34        let repositories = store
35            .list_code_repositories()
36            .await
37            .map_err(storage_api_error)?;
38        if repositories.is_empty() {
39            return Ok(false);
40        }
41        let repository_count = repositories.len();
42        let start = self.code_retention_cursor.fetch_add(1, Ordering::Relaxed) % repository_count;
43        let mut first_error = None;
44        let mut pending_repository = None;
45        for status in repositories
46            .into_iter()
47            .cycle()
48            .skip(start)
49            .take(repository_count)
50        {
51            let retention = match store
52                .code_scope_retention(status.repository_id.clone())
53                .await
54            {
55                Ok(retention) => retention,
56                Err(error) => {
57                    first_error.get_or_insert_with(|| storage_api_error(error));
58                    continue;
59                }
60            };
61            if retention.maintenance_pending && pending_repository.is_none() {
62                pending_repository = Some((
63                    status.repository_id,
64                    status.last_indexed_scope_id.unwrap_or_default(),
65                ));
66            }
67        }
68        let maintenance_active = if let Some((repository_id, active_scope)) = pending_repository {
69            match store
70                .prune_code_repository_scopes(crate::storage::CodeScopeRetentionRequest {
71                    repository_id,
72                    active_scope,
73                    retain_recent_successful_scopes: RETAIN_RECENT_CODE_SCOPES,
74                    repository_retention_cutoff_ms: None,
75                    repository_retention_cutoff_generation: None,
76                    repository_retention_initial_scope: None,
77                })
78                .await
79            {
80                Ok(pruned) => {
81                    pruned.pruned_scope_count > 0
82                        || pruned.retiring_job_count > 0
83                        || pruned.maintenance_pending
84                }
85                Err(error) => {
86                    first_error.get_or_insert_with(|| storage_api_error(error));
87                    false
88                }
89            }
90        } else {
91            false
92        };
93        first_error.map_or(Ok(repository_scan_pending || maintenance_active), Err)
94    }
95
96    /// Resets unfinished full index tasks for a registered repository.
97    pub async fn reset_code_repository_index_tasks(
98        &self,
99        repository: String,
100        context: RequestContext,
101    ) -> Result<CodeRepositoryIndexResetResponse, ApiError> {
102        let store = self.store().await.map_err(storage_api_error)?;
103        let status = required_code_repository(store.as_ref(), &repository).await?;
104        // Reclaim leases whose owner process has exited before applying the
105        // operator reset. This keeps reset idempotent while preserving the
106        // single-writer invariant for genuinely live workers.
107        recover_orphaned_code_index_task_leases(
108            &store,
109            now_millis(),
110            &self.runtime.process.windows_tasklist_command,
111        )
112        .await?;
113        let reset_tasks = store
114            .reset_code_index_tasks(status.repository_id.clone(), now_millis())
115            .await
116            .map_err(storage_api_error)?;
117        let active_task = store
118            .active_code_index_task(status.repository_id.clone())
119            .await
120            .map_err(storage_api_error)?;
121        let checkpoint =
122            code_status_checkpoint(store.as_ref(), &status, active_task.as_ref()).await?;
123        let retention = store
124            .code_scope_retention(status.repository_id.clone())
125            .await
126            .map_err(storage_api_error)?;
127        let graph_version = store
128            .current_graph_version()
129            .await
130            .map_err(storage_api_error)?;
131
132        Ok(CodeRepositoryIndexResetResponse {
133            metadata: ApiMetadata::graph_only(&context, graph_version),
134            status,
135            reset_task_count: reset_tasks.len(),
136            reset_tasks,
137            active_task,
138            checkpoint,
139            retention,
140        })
141    }
142
143    /// Reconciles expired or orphaned repository index leases before resident workers start.
144    pub async fn reconcile_startup_code_index_tasks(&self) -> Result<(), ApiError> {
145        let store = self.store().await.map_err(storage_api_error)?;
146        recover_orphaned_code_index_task_leases(
147            &store,
148            now_millis(),
149            &self.runtime.process.windows_tasklist_command,
150        )
151        .await
152        .map(|_| ())
153    }
154
155    pub(crate) async fn code_index_worker_status(
156        &self,
157        store: &std::sync::Arc<dyn crate::storage::KnowledgeStore>,
158    ) -> Result<CodeIndexWorkerStatus, ApiError> {
159        recover_orphaned_code_index_task_leases(
160            store,
161            now_millis(),
162            &self.runtime.process.windows_tasklist_command,
163        )
164        .await?;
165        self.read_only_code_index_worker_status(store).await
166    }
167
168    pub(crate) async fn read_only_code_index_worker_status(
169        &self,
170        store: &std::sync::Arc<dyn crate::storage::KnowledgeStore>,
171    ) -> Result<CodeIndexWorkerStatus, ApiError> {
172        let queue = store
173            .code_index_task_queue_status()
174            .await
175            .map_err(storage_api_error)?;
176
177        Ok(CodeIndexWorkerStatus::from_queue(
178            self.runtime.workers.code_index_max_in_flight,
179            queue,
180        ))
181    }
182}