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 CodeRepositorySetRefreshTaskClaimRequest, CodeRepositorySetRefreshTaskCompletion,
9 CodeRepositorySetRefreshTaskFailure, CodeRepositorySetRefreshTaskSeed,
10 },
11};
12
13use super::{
14 super::clock::now_millis,
15 errors::storage_api_error,
16 member_freshness::fact_version_scope_mismatch_reason,
17 status::{
18 persist_fact_version_member_replacements, refreshed_required_set_status,
19 required_set_status,
20 },
21};
22
23const REPOSITORY_SET_REFRESH_TASK_LEASE_MS: u64 = 10 * 60 * 1000;
24const REPOSITORY_SET_REFRESH_TASK_MAX_ATTEMPTS: u32 = 3;
25const REPOSITORY_SET_REFRESH_TASK_RETRY_BACKOFF_MS: u64 = 60_000;
26
27impl RelayKnowledgeService {
28 pub async fn refresh_code_repository_set(
30 &self,
31 set_alias: String,
32 context: RequestContext,
33 ) -> Result<CodeRepositorySetRefreshResponse, ApiError> {
34 let store = self.store().await.map_err(storage_api_error)?;
35 let (preflight_status, replacements) =
36 refreshed_required_set_status(&store, &set_alias).await?;
37 if let Some(reason) = preflight_status
38 .members
39 .iter()
40 .find_map(fact_version_scope_mismatch_reason)
41 {
42 return Err(ApiError::invalid_argument(format!(
43 "code repository set '{set_alias}' cannot refresh overlay: {reason}"
44 )));
45 }
46 persist_fact_version_member_replacements(&store, &set_alias, &replacements).await?;
47 let summary = store
48 .refresh_code_repository_set_overlay(set_alias.clone(), now_millis())
49 .await
50 .map_err(storage_api_error)?;
51 let status = required_set_status(&store, &set_alias).await?;
52 let graph_version = store
53 .current_graph_version()
54 .await
55 .map_err(storage_api_error)?;
56
57 Ok(CodeRepositorySetRefreshResponse {
58 metadata: ApiMetadata::graph_only(&context, graph_version),
59 status,
60 summary: Some(summary),
61 task: None,
62 })
63 }
64
65 pub async fn start_code_repository_set_refresh(
67 &self,
68 set_alias: String,
69 context: RequestContext,
70 ) -> Result<CodeRepositorySetRefreshResponse, ApiError> {
71 let store = self.store().await.map_err(storage_api_error)?;
72 let status = required_set_status(&store, &set_alias).await?;
73 let fingerprint = repository_set_refresh_fingerprint(&status);
74 let task = store
75 .queue_code_repository_set_refresh_task(CodeRepositorySetRefreshTaskSeed {
76 set_id: status.repository_set.set_id.clone(),
77 set_alias: status.repository_set.alias.clone(),
78 input_fingerprint: fingerprint,
79 now_ms: now_millis(),
80 })
81 .await
82 .map_err(storage_api_error)?;
83 let graph_version = store
84 .current_graph_version()
85 .await
86 .map_err(storage_api_error)?;
87
88 Ok(CodeRepositorySetRefreshResponse {
89 metadata: ApiMetadata::graph_only(&context, graph_version),
90 status,
91 summary: None,
92 task: Some(task),
93 })
94 }
95
96 pub async fn run_code_repository_set_refresh_task_once(
98 &self,
99 task_id: Option<String>,
100 context: RequestContext,
101 ) -> Result<Option<CodeRepositorySetRefreshTaskRecord>, ApiError> {
102 let store = self.store().await.map_err(storage_api_error)?;
103 let lease_owner = format!("code-repository-set-refresh-worker-{}", std::process::id());
104 let Some(task) = store
105 .claim_code_repository_set_refresh_task(CodeRepositorySetRefreshTaskClaimRequest {
106 task_id,
107 lease_owner: lease_owner.clone(),
108 lease_duration_ms: REPOSITORY_SET_REFRESH_TASK_LEASE_MS,
109 max_attempts: REPOSITORY_SET_REFRESH_TASK_MAX_ATTEMPTS,
110 now_ms: now_millis(),
111 })
112 .await
113 .map_err(storage_api_error)?
114 else {
115 return Ok(None);
116 };
117 let result = self
118 .refresh_code_repository_set(task.set_alias.clone(), context)
119 .await;
120 match result {
121 Ok(_) => store
122 .complete_code_repository_set_refresh_task(CodeRepositorySetRefreshTaskCompletion {
123 task_id: task.task_id,
124 lease_owner,
125 attempt_count: task.attempt_count,
126 now_ms: now_millis(),
127 })
128 .await
129 .map(Some)
130 .map_err(storage_api_error),
131 Err(error) => {
132 let _ = store
133 .fail_code_repository_set_refresh_task(CodeRepositorySetRefreshTaskFailure {
134 task_id: task.task_id,
135 lease_owner,
136 attempt_count: task.attempt_count,
137 error_kind: "repository_set_overlay".to_owned(),
138 error_message: error.message.clone(),
139 retry_backoff_ms: REPOSITORY_SET_REFRESH_TASK_RETRY_BACKOFF_MS,
140 max_attempts: REPOSITORY_SET_REFRESH_TASK_MAX_ATTEMPTS,
141 now_ms: now_millis(),
142 })
143 .await;
144 Err(error)
145 }
146 }
147 }
148}
149
150fn repository_set_refresh_fingerprint(status: &CodeRepositorySetStatus) -> String {
151 let mut parts = vec![status.repository_set.set_id.clone()];
152 parts.extend(status.members.iter().map(|member| {
153 format!(
154 "{}:{}:{}:{}:{}",
155 member.member.repository_id,
156 member.member.source_scope,
157 member.member.resolved_commit_sha,
158 member.tree_hash,
159 member.stale
160 )
161 }));
162 parts.join("|")
163}
164
165#[cfg(test)]
166#[path = "mod_tests.rs"]
167mod tests;