xbp-deploy 10.43.0

Service-centric declarative deploy engine for XBP.
Documentation
use std::path::{Path, PathBuf};
use std::sync::Arc;

use async_trait::async_trait;
use futures::future::join_all;
use xbp_oci::{OciRef, OciResolver};

use crate::context::DeployContext;
use crate::error::{DeployError, Result};
use crate::graph::{enforce_athena_operator_before_runtime, order_services, ServiceGraph};
use crate::lock::compute_plan_hash;
use crate::target_resolver::resolve_target;
use crate::types::*;
use crate::validator::validate_project_config;

#[async_trait]
pub trait DeployPlanner: Send + Sync {
    async fn plan(&self, ctx: &DeployContext) -> Result<DeployPlan>;
}

pub struct DefaultDeployPlanner {
    pub oci: Option<Arc<dyn OciResolver>>,
}

impl DefaultDeployPlanner {
    pub fn new(oci: Option<Arc<dyn OciResolver>>) -> Self {
        Self { oci }
    }
}

#[async_trait]
impl DeployPlanner for DefaultDeployPlanner {
    async fn plan(&self, ctx: &DeployContext) -> Result<DeployPlan> {
        validate_project_config(&ctx.config, &ctx.env)?;

        // Always re-resolve from label so group vs service priority is authoritative.
        let resolution = resolve_target(&ctx.config, &ctx.target.label(), &ctx.env)?;
        let target = resolution.target.clone();

        let graph = ServiceGraph::from_services(&ctx.config.services);
        let selected_names: Vec<String> = resolution.services.iter().map(|s| s.name.clone()).collect();
        let mut order = order_services(
            &selected_names,
            resolution.group_order.as_deref(),
            &graph,
        )?;
        if matches!(&target, DeployTarget::Group(g) if g == "athena") {
            enforce_athena_operator_before_runtime(&mut order);
        }

        // Operator priority: stable re-order operators first without breaking relative order.
        order.sort_by_key(|name| {
            resolution
                .services
                .iter()
                .find(|s| s.name == *name)
                .map(|s| {
                    provider_priority(
                        s.deploy
                            .as_ref()
                            .map(|d| d.provider.as_str())
                            .unwrap_or("kubernetes"),
                    )
                })
                .unwrap_or(9)
        });

        let mut service_plans = Vec::new();
        for name in &order {
            let svc = resolution
                .services
                .iter()
                .find(|s| s.name == *name)
                .ok_or_else(|| DeployError::Other(format!("missing service {name}")))?;
            service_plans.push(build_service_plan(&ctx.config, svc, &ctx.env, &ctx.flags)?);
        }

        // Parallel OCI digest resolution.
        if !ctx.flags.skip_oci {
            if let Some(oci) = &self.oci {
                let futures = service_plans.iter().map(|sp| {
                    let oci = Arc::clone(oci);
                    let image = sp.image.clone();
                    async move {
                        match image {
                            Some(img) => oci.resolve_digest(&img).await.ok(),
                            None => None,
                        }
                    }
                });
                let digests = join_all(futures).await;
                for (sp, digest) in service_plans.iter_mut().zip(digests) {
                    sp.digest = digest;
                }
            }
        }

        let mut oci_plan = OciPlan::default();
        let mut k8s_services = Vec::new();
        for sp in &service_plans {
            if let Some(r) = &sp.image_ref {
                oci_plan.images.insert(
                    sp.name.clone(),
                    OciImagePlan {
                        ref_str: r.clone(),
                        digest: sp.digest.clone(),
                    },
                );
            }
            if matches!(sp.provider.as_str(), "kubernetes" | "kubernetes-operator") {
                k8s_services.push(xbp_k8s::K8sServicePlan {
                    service: sp.name.clone(),
                    provider: sp.provider.clone(),
                    namespace: sp.deploy.namespace.clone().or_else(|| ctx.flags.namespace.clone()),
                    workload: sp.deploy.workload.clone(),
                    manifest_paths: sp.deploy.manifest_paths.clone(),
                    crds_path: sp.deploy.crds_path.clone(),
                    install_path: sp.deploy.install_path.clone(),
                    selector: sp.deploy.selector.clone(),
                    health: sp.deploy.health.clone(),
                });
            }
        }

        let mut plan = DeployPlan {
            target,
            env: ctx.env.clone(),
            project: ctx.config.project_name.clone(),
            project_version: ctx.config.version.clone(),
            git_sha: ctx.config.git_sha.clone(),
            services: service_plans,
            order,
            oci_plan,
            k8s_plan: K8sPlanView {
                context: ctx.flags.context.clone().or_else(|| {
                    ctx.config
                        .kubernetes
                        .as_ref()
                        .and_then(|k| k.default_context.clone())
                }),
                default_namespace: ctx.flags.namespace.clone().or_else(|| {
                    ctx.config
                        .kubernetes
                        .as_ref()
                        .and_then(|k| k.default_namespace.clone())
                }),
                services: k8s_services,
            },
            hash: String::new(),
        };
        plan.hash = compute_plan_hash(&plan);
        Ok(plan)
    }
}

fn provider_priority(provider: &str) -> u8 {
    match provider {
        "kubernetes-operator" => 0,
        "kubernetes" => 1,
        "worker" => 2,
        _ => 9,
    }
}

fn build_service_plan(
    config: &ProjectConfig,
    svc: &ServiceConfigView,
    env: &str,
    flags: &crate::context::DeployFlags,
) -> Result<ServicePlan> {
    let deploy = svc
        .deploy
        .as_ref()
        .ok_or_else(|| DeployError::Validation(format!("service `{}` has no deploy", svc.name)))?;
    let env_cfg = deploy.envs.get(env).ok_or_else(|| {
        DeployError::Validation(format!("service `{}` missing deploy.envs.{env}", svc.name))
    })?;
    let provider = deploy.provider.clone();
    let version = svc
        .version
        .clone()
        .unwrap_or_else(|| config.version.clone());

    let (image, image_ref) = resolve_image(svc, &version, &config.version);
    let namespace = flags
        .namespace
        .clone()
        .or_else(|| env_cfg.namespace.clone())
        .or_else(|| {
            config
                .kubernetes
                .as_ref()
                .and_then(|k| k.default_namespace.clone())
        });

    let mut manifest_paths = Vec::new();
    let mut crds_path = None;
    let mut install_path = None;
    let mut workload = None;
    let mut selector = None;
    if let Some(k8s) = &env_cfg.kubernetes {
        for candidate in [&k8s.manifests_base, &k8s.manifests_overlay]
            .into_iter()
            .flatten()
        {
            manifest_paths.push(resolve_path(
                &config.project_root,
                svc.root_directory.as_deref(),
                candidate,
            ));
        }
        workload = k8s.workload.clone();
        selector = k8s.selector.clone();
        crds_path = k8s.crds_path.as_ref().map(|p| {
            resolve_path(&config.project_root, svc.root_directory.as_deref(), p)
        });
        install_path = k8s.install_path.as_ref().map(|p| {
            resolve_path(&config.project_root, svc.root_directory.as_deref(), p)
        });
    }

    let mut actions = Vec::new();
    match provider.as_str() {
        "kubernetes-operator" => {
            if crds_path.is_some() {
                actions.push("apply CRDs".into());
            }
            if install_path.is_some() {
                actions.push("apply operator install".into());
            }
            if selector.is_some() {
                actions.push("wait operator selector".into());
            }
        }
        "kubernetes" => {
            if image_ref.is_some() {
                actions.push("verify OCI image".into());
            }
            if !manifest_paths.is_empty() {
                actions.push(format!("apply {} manifest path(s)", manifest_paths.len()));
            }
            if workload.is_some() {
                actions.push("wait rollout".into());
            }
            if !env_cfg.health.is_empty() {
                actions.push(format!("probe {} health URL(s)", env_cfg.health.len()));
            }
        }
        other => actions.push(format!("provider `{other}`")),
    }

    Ok(ServicePlan {
        name: svc.name.clone(),
        provider,
        version,
        image,
        image_ref,
        digest: None,
        deploy: ServiceDeployPlan {
            namespace,
            workload,
            health: env_cfg.health.clone(),
            manifest_paths,
            crds_path,
            install_path,
            selector,
            actions,
        },
    })
}

fn resolve_image(
    svc: &ServiceConfigView,
    service_version: &str,
    project_version: &str,
) -> (Option<OciRef>, Option<String>) {
    let Some(oci) = &svc.oci else {
        return (None, None);
    };
    let tag = if let Some(t) = oci.tag.as_deref().filter(|s| !s.is_empty()) {
        t.to_string()
    } else {
        match oci.tag_from.as_deref().unwrap_or("project_version") {
            "service_version" => service_version.to_string(),
            "git_sha" => "git".into(),
            _ => project_version.to_string(),
        }
    };
    let ref_str = format!("{}:{tag}", oci.image.trim());
    match xbp_oci::parse_image_ref(&ref_str) {
        Ok(r) => (Some(r), Some(ref_str)),
        Err(_) => (None, Some(ref_str)),
    }
}

fn resolve_path(project_root: &Path, service_root: Option<&str>, relative: &str) -> PathBuf {
    let rel = PathBuf::from(relative);
    if rel.is_absolute() {
        return rel;
    }
    if let Some(root) = service_root {
        let under = project_root.join(root).join(&rel);
        if under.exists() {
            return under;
        }
    }
    project_root.join(rel)
}