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        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    /// Rebuilds cross-repository import/module overlay edges.
29    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    /// Queues a repository-set overlay refresh task.
66    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    /// Runs one queued repository-set overlay refresh task under a lease.
97    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;