relay_knowledge/application/code_repository/repository_set/refresh/
mod.rs1use crate::{
4 api::{ApiError, ApiMetadata, CodeRepositorySetRefreshResponse, RequestContext},
5 application::service::RelayKnowledgeService,
6 domain::{CodeRepositorySetRefreshTaskRecord, CodeRepositorySetStatus},
7 storage::{
8 CodeRepositorySetMemberSeed, CodeRepositorySetRefreshPublication,
9 CodeRepositorySetRefreshTaskClaimRequest, CodeRepositorySetRefreshTaskCompletion,
10 CodeRepositorySetRefreshTaskFailure, CodeRepositorySetRefreshTaskSeed,
11 },
12};
13
14use super::{
15 super::clock::now_millis,
16 errors::storage_api_error,
17 member_freshness::fact_version_scope_mismatch_reason,
18 status::{refreshed_required_set_status, required_set_status},
19};
20
21const REPOSITORY_SET_REFRESH_TASK_LEASE_MS: u64 = 10 * 60 * 1000;
22const REPOSITORY_SET_REFRESH_TASK_MAX_ATTEMPTS: u32 = 3;
23const REPOSITORY_SET_REFRESH_TASK_RETRY_BACKOFF_MS: u64 = 60_000;
24
25struct RefreshTaskExecution {
26 task: CodeRepositorySetRefreshTaskRecord,
27 response: CodeRepositorySetRefreshResponse,
28}
29
30impl RelayKnowledgeService {
31 pub async fn refresh_code_repository_set(
33 &self,
34 set_alias: String,
35 context: RequestContext,
36 ) -> Result<CodeRepositorySetRefreshResponse, ApiError> {
37 let queued = self
38 .start_code_repository_set_refresh(set_alias, context.clone())
39 .await?;
40 let task_id = queued.task.as_ref().map(|task| task.task_id.clone());
41 match self
42 .execute_code_repository_set_refresh_task(task_id, context)
43 .await?
44 {
45 Some(execution) => Ok(execution.response),
46 None => Ok(queued),
47 }
48 }
49
50 pub async fn start_code_repository_set_refresh(
52 &self,
53 set_alias: String,
54 context: RequestContext,
55 ) -> Result<CodeRepositorySetRefreshResponse, ApiError> {
56 let store = self.store().await.map_err(storage_api_error)?;
57 let status = required_set_status(&store, &set_alias).await?;
58 let fingerprint = repository_set_refresh_fingerprint(&status);
59 let task = store
60 .queue_code_repository_set_refresh_task(CodeRepositorySetRefreshTaskSeed {
61 set_id: status.repository_set.set_id.clone(),
62 set_alias: status.repository_set.alias.clone(),
63 input_fingerprint: fingerprint,
64 now_ms: now_millis(),
65 })
66 .await
67 .map_err(storage_api_error)?;
68 let graph_version = store
69 .current_graph_version()
70 .await
71 .map_err(storage_api_error)?;
72
73 Ok(CodeRepositorySetRefreshResponse {
74 metadata: ApiMetadata::graph_only(&context, graph_version),
75 status,
76 summary: None,
77 task: Some(task),
78 })
79 }
80
81 pub async fn run_code_repository_set_refresh_task_once(
83 &self,
84 task_id: Option<String>,
85 context: RequestContext,
86 ) -> Result<Option<CodeRepositorySetRefreshTaskRecord>, ApiError> {
87 self.execute_code_repository_set_refresh_task(task_id, context)
88 .await
89 .map(|execution| execution.map(|execution| execution.task))
90 }
91
92 async fn execute_code_repository_set_refresh_task(
93 &self,
94 task_id: Option<String>,
95 context: RequestContext,
96 ) -> Result<Option<RefreshTaskExecution>, ApiError> {
97 let store = self.store().await.map_err(storage_api_error)?;
98 let lease_owner = format!("code-repository-set-refresh-worker-{}", std::process::id());
99 let Some(task) = store
100 .claim_code_repository_set_refresh_task(CodeRepositorySetRefreshTaskClaimRequest {
101 task_id,
102 lease_owner: lease_owner.clone(),
103 lease_duration_ms: REPOSITORY_SET_REFRESH_TASK_LEASE_MS,
104 max_attempts: REPOSITORY_SET_REFRESH_TASK_MAX_ATTEMPTS,
105 now_ms: now_millis(),
106 })
107 .await
108 .map_err(storage_api_error)?
109 else {
110 return Ok(None);
111 };
112 let result = self
113 .build_and_publish_code_repository_set_refresh(&task, &lease_owner, context)
114 .await;
115 match result {
116 Ok(mut response) => {
117 let completed = store
118 .complete_code_repository_set_refresh_task(
119 CodeRepositorySetRefreshTaskCompletion {
120 task_id: task.task_id,
121 lease_owner,
122 attempt_count: task.attempt_count,
123 now_ms: now_millis(),
124 },
125 )
126 .await
127 .map_err(storage_api_error)?;
128 response.task = Some(completed.clone());
129 Ok(Some(RefreshTaskExecution {
130 task: completed,
131 response,
132 }))
133 }
134 Err(error) => {
135 let _ = store
136 .fail_code_repository_set_refresh_task(CodeRepositorySetRefreshTaskFailure {
137 task_id: task.task_id,
138 lease_owner,
139 attempt_count: task.attempt_count,
140 error_kind: "repository_set_overlay".to_owned(),
141 error_message: error.message.clone(),
142 retry_backoff_ms: REPOSITORY_SET_REFRESH_TASK_RETRY_BACKOFF_MS,
143 max_attempts: REPOSITORY_SET_REFRESH_TASK_MAX_ATTEMPTS,
144 now_ms: now_millis(),
145 })
146 .await;
147 Err(error)
148 }
149 }
150 }
151
152 async fn build_and_publish_code_repository_set_refresh(
153 &self,
154 task: &CodeRepositorySetRefreshTaskRecord,
155 lease_owner: &str,
156 context: RequestContext,
157 ) -> Result<CodeRepositorySetRefreshResponse, ApiError> {
158 let store = self.store().await.map_err(storage_api_error)?;
159 let (preflight_status, replacements) =
160 refreshed_required_set_status(&store, &task.set_alias).await?;
161 if let Some(reason) = preflight_status
162 .members
163 .iter()
164 .find_map(fact_version_scope_mismatch_reason)
165 {
166 return Err(ApiError::invalid_argument(format!(
167 "code repository set '{}' cannot refresh overlay: {reason}",
168 task.set_alias
169 )));
170 }
171 let summary = store
172 .refresh_code_repository_set_overlay(
173 task.set_alias.clone(),
174 CodeRepositorySetRefreshPublication {
175 task_id: task.task_id.clone(),
176 set_id: task.set_id.clone(),
177 lease_owner: lease_owner.to_owned(),
178 attempt_count: task.attempt_count,
179 member_replacements: replacements
180 .into_iter()
181 .map(|member| CodeRepositorySetMemberSeed {
182 set_alias: task.set_alias.clone(),
183 repository_id: member.repository_id,
184 repository_alias: member.repository_alias,
185 ref_selector: member.ref_selector,
186 resolved_commit_sha: member.resolved_commit_sha,
187 source_scope: member.source_scope,
188 path_filters: member.path_filters,
189 language_filters: member.language_filters,
190 priority: member.priority,
191 })
192 .collect(),
193 },
194 )
195 .await
196 .map_err(storage_api_error)?;
197 let status = required_set_status(&store, &task.set_alias).await?;
198 let graph_version = store
199 .current_graph_version()
200 .await
201 .map_err(storage_api_error)?;
202
203 Ok(CodeRepositorySetRefreshResponse {
204 metadata: ApiMetadata::graph_only(&context, graph_version),
205 status,
206 summary: Some(summary),
207 task: None,
208 })
209 }
210}
211
212fn repository_set_refresh_fingerprint(status: &CodeRepositorySetStatus) -> String {
213 let mut parts = vec![status.repository_set.set_id.clone()];
214 parts.extend(status.members.iter().map(|member| {
215 format!(
216 "{}:{}:{}:{}:{}",
217 member.member.repository_id,
218 member.member.source_scope,
219 member.member.resolved_commit_sha,
220 member.tree_hash,
221 member.stale
222 )
223 }));
224 parts.join("|")
225}
226
227#[cfg(test)]
228#[path = "mod_tests.rs"]
229mod tests;