Skip to main content

relay_knowledge/application/code_repository/repository_set/refresh/
mod.rs

1//! Coordinates synchronous and durable repository-set overlay refresh.
2
3use 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    /// Queues and, when available, synchronously drains one overlay refresh task.
32    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    /// Queues a repository-set overlay refresh task.
51    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    /// Runs one queued repository-set overlay refresh task under a lease.
82    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;