Skip to main content

relay_knowledge/application/code_repository/
tasks.rs

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