Skip to main content

relay_knowledge/application/code_repository/indexing/
tasks.rs

1use crate::{
2    api::{
3        ApiError, ApiMetadata, CodeIndexWorkerStatus, CodeRepositoryIndexResetResponse,
4        RequestContext,
5    },
6    application::service::RelayKnowledgeService,
7};
8
9use super::super::{
10    clock::now_millis,
11    errors::storage_api_error,
12    repository::{code_status_checkpoint, required_code_repository},
13};
14use super::task::recover_orphaned_code_index_task_leases;
15
16impl RelayKnowledgeService {
17    /// Resets unfinished full index tasks for a registered repository.
18    pub async fn reset_code_repository_index_tasks(
19        &self,
20        repository: String,
21        context: RequestContext,
22    ) -> Result<CodeRepositoryIndexResetResponse, ApiError> {
23        let store = self.store().await.map_err(storage_api_error)?;
24        let status = required_code_repository(&store, &repository).await?;
25        // Reclaim leases whose owner process has exited before applying the
26        // operator reset. This keeps reset idempotent while preserving the
27        // single-writer invariant for genuinely live workers.
28        recover_orphaned_code_index_task_leases(
29            &store,
30            now_millis(),
31            &self.runtime.process.windows_tasklist_command,
32        )
33        .await?;
34        let reset_tasks = store
35            .reset_code_index_tasks(status.repository_id.clone(), now_millis())
36            .await
37            .map_err(storage_api_error)?;
38        let active_task = store
39            .active_code_index_task(status.repository_id.clone())
40            .await
41            .map_err(storage_api_error)?;
42        let checkpoint = code_status_checkpoint(&store, &status, active_task.as_ref()).await?;
43        let retention = store
44            .code_scope_retention(status.repository_id.clone())
45            .await
46            .map_err(storage_api_error)?;
47        let graph_version = store
48            .current_graph_version()
49            .await
50            .map_err(storage_api_error)?;
51
52        Ok(CodeRepositoryIndexResetResponse {
53            metadata: ApiMetadata::graph_only(&context, graph_version),
54            status,
55            reset_task_count: reset_tasks.len(),
56            reset_tasks,
57            active_task,
58            checkpoint,
59            retention,
60        })
61    }
62
63    /// Reconciles expired or orphaned repository index leases before resident workers start.
64    pub async fn reconcile_startup_code_index_tasks(&self) -> Result<(), ApiError> {
65        let store = self.store().await.map_err(storage_api_error)?;
66        recover_orphaned_code_index_task_leases(
67            &store,
68            now_millis(),
69            &self.runtime.process.windows_tasklist_command,
70        )
71        .await
72        .map(|_| ())
73    }
74
75    pub(crate) async fn code_index_worker_status(
76        &self,
77        store: &std::sync::Arc<dyn crate::storage::KnowledgeStore>,
78    ) -> Result<CodeIndexWorkerStatus, ApiError> {
79        recover_orphaned_code_index_task_leases(
80            store,
81            now_millis(),
82            &self.runtime.process.windows_tasklist_command,
83        )
84        .await?;
85        self.read_only_code_index_worker_status(store).await
86    }
87
88    pub(crate) async fn read_only_code_index_worker_status(
89        &self,
90        store: &std::sync::Arc<dyn crate::storage::KnowledgeStore>,
91    ) -> Result<CodeIndexWorkerStatus, ApiError> {
92        let queue = store
93            .code_index_task_queue_status()
94            .await
95            .map_err(storage_api_error)?;
96
97        Ok(CodeIndexWorkerStatus::from_queue(
98            self.runtime.workers.code_index_max_in_flight,
99            queue,
100        ))
101    }
102}