Skip to main content

xbp_deploy/
planner.rs

1use std::path::{Path, PathBuf};
2use std::sync::Arc;
3
4use async_trait::async_trait;
5use futures::future::join_all;
6use xbp_oci::{OciRef, OciResolver};
7
8use crate::context::DeployContext;
9use crate::error::{DeployError, Result};
10use crate::graph::{enforce_athena_operator_before_runtime, order_services, ServiceGraph};
11use crate::lock::compute_plan_hash;
12use crate::target_resolver::resolve_target;
13use crate::types::*;
14use crate::validator::validate_project_config;
15
16#[async_trait]
17pub trait DeployPlanner: Send + Sync {
18    async fn plan(&self, ctx: &DeployContext) -> Result<DeployPlan>;
19}
20
21pub struct DefaultDeployPlanner {
22    pub oci: Option<Arc<dyn OciResolver>>,
23}
24
25impl DefaultDeployPlanner {
26    pub fn new(oci: Option<Arc<dyn OciResolver>>) -> Self {
27        Self { oci }
28    }
29}
30
31#[async_trait]
32impl DeployPlanner for DefaultDeployPlanner {
33    async fn plan(&self, ctx: &DeployContext) -> Result<DeployPlan> {
34        validate_project_config(&ctx.config, &ctx.env)?;
35
36        // Always re-resolve from label so group vs service priority is authoritative.
37        let resolution = resolve_target(&ctx.config, &ctx.target.label(), &ctx.env)?;
38        let target = resolution.target.clone();
39
40        let graph = ServiceGraph::from_services(&ctx.config.services);
41        let selected_names: Vec<String> = resolution.services.iter().map(|s| s.name.clone()).collect();
42        let mut order = order_services(
43            &selected_names,
44            resolution.group_order.as_deref(),
45            &graph,
46        )?;
47        if matches!(&target, DeployTarget::Group(g) if g == "athena") {
48            enforce_athena_operator_before_runtime(&mut order);
49        }
50
51        // Operator priority: stable re-order operators first without breaking relative order.
52        order.sort_by_key(|name| {
53            resolution
54                .services
55                .iter()
56                .find(|s| s.name == *name)
57                .map(|s| {
58                    provider_priority(
59                        s.deploy
60                            .as_ref()
61                            .map(|d| d.provider.as_str())
62                            .unwrap_or("kubernetes"),
63                    )
64                })
65                .unwrap_or(9)
66        });
67
68        let mut service_plans = Vec::new();
69        for name in &order {
70            let svc = resolution
71                .services
72                .iter()
73                .find(|s| s.name == *name)
74                .ok_or_else(|| DeployError::Other(format!("missing service {name}")))?;
75            service_plans.push(build_service_plan(&ctx.config, svc, &ctx.env, &ctx.flags)?);
76        }
77
78        // Parallel OCI digest resolution.
79        if !ctx.flags.skip_oci {
80            if let Some(oci) = &self.oci {
81                let futures = service_plans.iter().map(|sp| {
82                    let oci = Arc::clone(oci);
83                    let image = sp.image.clone();
84                    async move {
85                        match image {
86                            Some(img) => oci.resolve_digest(&img).await.ok(),
87                            None => None,
88                        }
89                    }
90                });
91                let digests = join_all(futures).await;
92                for (sp, digest) in service_plans.iter_mut().zip(digests) {
93                    sp.digest = digest;
94                }
95            }
96        }
97
98        let mut oci_plan = OciPlan::default();
99        let mut k8s_services = Vec::new();
100        for sp in &service_plans {
101            if let Some(r) = &sp.image_ref {
102                oci_plan.images.insert(
103                    sp.name.clone(),
104                    OciImagePlan {
105                        ref_str: r.clone(),
106                        digest: sp.digest.clone(),
107                    },
108                );
109            }
110            if matches!(sp.provider.as_str(), "kubernetes" | "kubernetes-operator") {
111                k8s_services.push(xbp_k8s::K8sServicePlan {
112                    service: sp.name.clone(),
113                    provider: sp.provider.clone(),
114                    namespace: sp.deploy.namespace.clone().or_else(|| ctx.flags.namespace.clone()),
115                    workload: sp.deploy.workload.clone(),
116                    manifest_paths: sp.deploy.manifest_paths.clone(),
117                    crds_path: sp.deploy.crds_path.clone(),
118                    install_path: sp.deploy.install_path.clone(),
119                    selector: sp.deploy.selector.clone(),
120                    health: sp.deploy.health.clone(),
121                });
122            }
123        }
124
125        let mut plan = DeployPlan {
126            target,
127            env: ctx.env.clone(),
128            project: ctx.config.project_name.clone(),
129            project_version: ctx.config.version.clone(),
130            git_sha: ctx.config.git_sha.clone(),
131            services: service_plans,
132            order,
133            oci_plan,
134            k8s_plan: K8sPlanView {
135                context: ctx.flags.context.clone().or_else(|| {
136                    ctx.config
137                        .kubernetes
138                        .as_ref()
139                        .and_then(|k| k.default_context.clone())
140                }),
141                default_namespace: ctx.flags.namespace.clone().or_else(|| {
142                    ctx.config
143                        .kubernetes
144                        .as_ref()
145                        .and_then(|k| k.default_namespace.clone())
146                }),
147                services: k8s_services,
148            },
149            hash: String::new(),
150        };
151        plan.hash = compute_plan_hash(&plan);
152        Ok(plan)
153    }
154}
155
156fn provider_priority(provider: &str) -> u8 {
157    match provider {
158        "kubernetes-operator" => 0,
159        "kubernetes" => 1,
160        "worker" => 2,
161        _ => 9,
162    }
163}
164
165fn build_service_plan(
166    config: &ProjectConfig,
167    svc: &ServiceConfigView,
168    env: &str,
169    flags: &crate::context::DeployFlags,
170) -> Result<ServicePlan> {
171    let deploy = svc
172        .deploy
173        .as_ref()
174        .ok_or_else(|| DeployError::Validation(format!("service `{}` has no deploy", svc.name)))?;
175    let env_cfg = deploy.envs.get(env).ok_or_else(|| {
176        DeployError::Validation(format!("service `{}` missing deploy.envs.{env}", svc.name))
177    })?;
178    let provider = deploy.provider.clone();
179    let version = svc
180        .version
181        .clone()
182        .unwrap_or_else(|| config.version.clone());
183
184    let (image, image_ref) = resolve_image(svc, &version, &config.version);
185    let namespace = flags
186        .namespace
187        .clone()
188        .or_else(|| env_cfg.namespace.clone())
189        .or_else(|| {
190            config
191                .kubernetes
192                .as_ref()
193                .and_then(|k| k.default_namespace.clone())
194        });
195
196    let mut manifest_paths = Vec::new();
197    let mut crds_path = None;
198    let mut install_path = None;
199    let mut workload = None;
200    let mut selector = None;
201    if let Some(k8s) = &env_cfg.kubernetes {
202        for candidate in [&k8s.manifests_base, &k8s.manifests_overlay]
203            .into_iter()
204            .flatten()
205        {
206            manifest_paths.push(resolve_path(
207                &config.project_root,
208                svc.root_directory.as_deref(),
209                candidate,
210            ));
211        }
212        workload = k8s.workload.clone();
213        selector = k8s.selector.clone();
214        crds_path = k8s.crds_path.as_ref().map(|p| {
215            resolve_path(&config.project_root, svc.root_directory.as_deref(), p)
216        });
217        install_path = k8s.install_path.as_ref().map(|p| {
218            resolve_path(&config.project_root, svc.root_directory.as_deref(), p)
219        });
220    }
221
222    let mut actions = Vec::new();
223    match provider.as_str() {
224        "kubernetes-operator" => {
225            if crds_path.is_some() {
226                actions.push("apply CRDs".into());
227            }
228            if install_path.is_some() {
229                actions.push("apply operator install".into());
230            }
231            if selector.is_some() {
232                actions.push("wait operator selector".into());
233            }
234        }
235        "kubernetes" => {
236            if image_ref.is_some() {
237                actions.push("verify OCI image".into());
238            }
239            if !manifest_paths.is_empty() {
240                actions.push(format!("apply {} manifest path(s)", manifest_paths.len()));
241            }
242            if workload.is_some() {
243                actions.push("wait rollout".into());
244            }
245            if !env_cfg.health.is_empty() {
246                actions.push(format!("probe {} health URL(s)", env_cfg.health.len()));
247            }
248        }
249        other => actions.push(format!("provider `{other}`")),
250    }
251
252    Ok(ServicePlan {
253        name: svc.name.clone(),
254        provider,
255        version,
256        image,
257        image_ref,
258        digest: None,
259        deploy: ServiceDeployPlan {
260            namespace,
261            workload,
262            health: env_cfg.health.clone(),
263            manifest_paths,
264            crds_path,
265            install_path,
266            selector,
267            actions,
268        },
269    })
270}
271
272fn resolve_image(
273    svc: &ServiceConfigView,
274    service_version: &str,
275    project_version: &str,
276) -> (Option<OciRef>, Option<String>) {
277    let Some(oci) = &svc.oci else {
278        return (None, None);
279    };
280    let tag = if let Some(t) = oci.tag.as_deref().filter(|s| !s.is_empty()) {
281        t.to_string()
282    } else {
283        match oci.tag_from.as_deref().unwrap_or("project_version") {
284            "service_version" => service_version.to_string(),
285            "git_sha" => "git".into(),
286            _ => project_version.to_string(),
287        }
288    };
289    let ref_str = format!("{}:{tag}", oci.image.trim());
290    match xbp_oci::parse_image_ref(&ref_str) {
291        Ok(r) => (Some(r), Some(ref_str)),
292        Err(_) => (None, Some(ref_str)),
293    }
294}
295
296fn resolve_path(project_root: &Path, service_root: Option<&str>, relative: &str) -> PathBuf {
297    let rel = PathBuf::from(relative);
298    if rel.is_absolute() {
299        return rel;
300    }
301    if let Some(root) = service_root {
302        let under = project_root.join(root).join(&rel);
303        if under.exists() {
304            return under;
305        }
306    }
307    project_root.join(rel)
308}