relay_knowledge/application/code_repository/indexing/
tasks.rs1use 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 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 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 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 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}