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)?;
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);
}
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)?);
}
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)
}