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::providers::{
destination_to_provider, is_all_destinations, is_cloudflare_provider, is_kubernetes_provider,
is_railway_provider, normalize_destination, provider_priority, provider_to_destination,
DESTINATION_ALL,
};
use crate::target_resolver::resolve_target;
use crate::types::*;
use crate::validator::validate_project_config_for_target;
#[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_for_target(&ctx.config, &ctx.env, Some(&ctx.target))?;
let resolution = resolve_target(&ctx.config, &ctx.target.label(), &ctx.env)?;
let target = resolution.target.clone();
if target != ctx.target {
validate_project_config_for_target(&ctx.config, &ctx.env, Some(&target))?;
}
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();
let mut expanded_order = 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}")))?;
let plans = build_service_plans(&ctx.config, svc, &ctx.env, &ctx.flags)?;
for plan in plans {
let step_key = match &plan.destination {
Some(d) => format!("{}@{}", plan.name, d),
None => plan.name.clone(),
};
expanded_order.push(step_key);
service_plans.push(plan);
}
}
let order = expanded_order;
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 is_kubernetes_provider(&sp.provider) {
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: {
let flag_or_project = ctx.flags.namespace.clone().or_else(|| {
ctx.config
.kubernetes
.as_ref()
.and_then(|k| k.default_namespace.clone())
});
let mut service_ns = k8s_services
.iter()
.filter_map(|s| s.namespace.as_ref())
.map(|s| s.trim())
.filter(|s| !s.is_empty());
if let Some(first) = service_ns.next() {
if service_ns.all(|n| n == first) {
Some(first.to_string())
} else {
flag_or_project
}
} else {
flag_or_project
}
},
services: k8s_services,
},
hash: String::new(),
};
plan.hash = compute_plan_hash(&plan);
Ok(plan)
}
}
struct DestinationSpec {
destination: String,
provider: String,
worker: Option<String>,
rollout: Option<String>,
}
fn configured_enabled_destinations(deploy: &ServiceDeployView) -> Vec<String> {
let mut scored: Vec<(u8, String)> = deploy
.destinations
.iter()
.filter(|(_, overlay)| overlay.enabled != Some(false))
.map(|(key, overlay)| {
let dest = normalize_destination(key);
let provider = overlay
.provider
.as_deref()
.filter(|s| !s.trim().is_empty())
.map(str::to_string)
.unwrap_or_else(|| destination_to_provider(&dest));
(provider_priority(&provider), dest)
})
.collect();
scored.sort_by(|a, b| a.0.cmp(&b.0).then_with(|| a.1.cmp(&b.1)));
let mut out = Vec::new();
for (_, dest) in scored {
if !out.iter().any(|d| d == &dest) {
out.push(dest);
}
}
out
}
fn resolve_destination_specs(
deploy: &ServiceDeployView,
service_name: &str,
flags: &crate::context::DeployFlags,
) -> Result<Vec<DestinationSpec>> {
let explicit_cli = !flags.destinations.is_empty();
let expand_all = explicit_cli && is_all_destinations(&flags.destinations);
let selected: Vec<String> = if expand_all {
let configured = configured_enabled_destinations(deploy);
if configured.is_empty() {
vec![provider_to_destination(&deploy.provider)]
} else {
configured
}
} else if explicit_cli {
flags
.destinations
.iter()
.map(|d| normalize_destination(d))
.filter(|d| d != DESTINATION_ALL)
.collect()
} else if !deploy.destinations.is_empty() {
let configured = configured_enabled_destinations(deploy);
if configured.is_empty() {
vec![provider_to_destination(&deploy.provider)]
} else {
configured
}
} else {
vec![provider_to_destination(&deploy.provider)]
};
let mut specs = Vec::new();
for dest in selected {
let overlay = deploy.destinations.get(&dest).or_else(|| {
deploy.destinations.iter().find_map(|(k, v)| {
if normalize_destination(k) == dest {
Some(v)
} else {
None
}
})
});
if let Some(overlay) = overlay {
if overlay.enabled == Some(false) && !explicit_cli {
continue;
}
let provider = overlay
.provider
.clone()
.filter(|s| !s.trim().is_empty())
.unwrap_or_else(|| destination_to_provider(&dest));
specs.push(DestinationSpec {
destination: dest.clone(),
provider,
worker: overlay
.worker
.clone()
.or_else(|| deploy.worker.clone())
.filter(|s| !s.trim().is_empty())
.or_else(|| Some(service_name.to_string())),
rollout: overlay.rollout.clone().or_else(|| deploy.rollout.clone()),
});
} else {
let provider = if !explicit_cli {
deploy.provider.clone()
} else {
destination_to_provider(&dest)
};
let use_service_worker = is_cloudflare_provider(&provider)
|| (!explicit_cli && is_cloudflare_provider(&deploy.provider));
specs.push(DestinationSpec {
destination: dest,
provider: if !explicit_cli {
deploy.provider.clone()
} else {
provider
},
worker: if use_service_worker {
deploy
.worker
.clone()
.filter(|s| !s.trim().is_empty())
.or_else(|| Some(service_name.to_string()))
} else {
deploy.worker.clone().filter(|s| !s.trim().is_empty())
},
rollout: deploy.rollout.clone(),
});
}
}
specs.sort_by(|a, b| {
provider_priority(&a.provider)
.cmp(&provider_priority(&b.provider))
.then_with(|| a.destination.cmp(&b.destination))
});
if specs.is_empty() {
return Err(DeployError::Validation(format!(
"service `{service_name}` has no deploy destinations after filters \
(configure services[].deploy.destinations or pass --to cloudflare|kubernetes|all)"
)));
}
Ok(specs)
}
fn build_service_plans(
config: &ProjectConfig,
svc: &ServiceConfigView,
env: &str,
flags: &crate::context::DeployFlags,
) -> Result<Vec<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 specs = resolve_destination_specs(deploy, &svc.name, flags)?;
let version = svc
.version
.clone()
.unwrap_or_else(|| config.version.clone());
let (image, image_ref) = resolve_image(svc, &version, &config.version);
let needs_strict_env = specs.iter().any(|s| is_kubernetes_provider(&s.provider));
let composed = crate::compose::compose_service_runtime_with_opts(
&config.project_root,
svc,
env,
env_cfg,
needs_strict_env,
)?;
if needs_strict_env && !composed.missing_required.is_empty() {
return Err(DeployError::Validation(format!(
"service `{}` missing required env for deploy.envs.{}: {} (set in .env, services.environment, or deploy.envs.{}.env)",
svc.name,
env,
composed.missing_required.join(", "),
env
)));
}
let mut plans = Vec::new();
for spec in specs {
plans.push(build_one_service_plan(
config,
svc,
env_cfg,
&version,
image.clone(),
image_ref.clone(),
flags,
spec,
&composed,
)?);
}
Ok(plans)
}
fn build_one_service_plan(
config: &ProjectConfig,
svc: &ServiceConfigView,
env_cfg: &ServiceDeployEnvView,
version: &str,
image: Option<OciRef>,
image_ref: Option<String>,
flags: &crate::context::DeployFlags,
spec: DestinationSpec,
composed: &crate::compose::ComposedRuntime,
) -> Result<ServicePlan> {
let provider = spec.provider;
let destination = spec.destination;
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 k8s_service = None;
let mut selector = None;
if is_kubernetes_provider(&provider) {
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();
k8s_service = k8s
.service
.as_ref()
.map(|s| s.trim().to_string())
.filter(|s| !s.is_empty());
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 worker_app = if is_cloudflare_provider(&provider) {
Some(
spec.worker
.clone()
.filter(|s| !s.trim().is_empty())
.unwrap_or_else(|| svc.name.clone()),
)
} else {
None
};
let rollout = if is_cloudflare_provider(&provider) {
spec.rollout.clone()
} else {
None
};
let mut actions = Vec::new();
let dockerfile = svc.oci.as_ref().and_then(|o| o.dockerfile.clone());
let build_context = svc.oci.as_ref().and_then(|o| o.context.clone());
let platforms = svc
.oci
.as_ref()
.map(|o| o.platforms.clone())
.unwrap_or_default();
actions.push(format!("destination `{destination}`"));
if is_kubernetes_provider(&provider) {
if dockerfile.is_some() && image_ref.is_some() && !flags.skip_build {
if flags.local_image {
actions.push(
"build local OCI image (docker load — no registry push)".into(),
);
} else {
actions.push("build + push OCI image (before cluster apply)".into());
}
} else if image_ref.is_some() {
actions.push("verify OCI image".into());
}
if destination == "kubernetes-operator"
|| provider.to_ascii_lowercase().contains("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());
}
} else {
if !manifest_paths.is_empty() {
actions.push(format!("apply {} manifest path(s)", manifest_paths.len()));
} else if workload.is_some() {
if flags.bootstrap_workload {
actions.push(
"bootstrap minimal Deployment/Service (no manifests configured)".into(),
);
} else {
actions.push(
"error risk: no kubernetes manifests — use --bootstrap-workload or set kubernetes.manifests".into(),
);
}
}
if !composed.env.is_empty() {
actions.push(format!(
"compose {} runtime env key(s) from .env / services.environment / deploy.envs",
composed.env.len()
));
}
if !composed.config_mounts.is_empty() {
actions.push(format!(
"mount {} config file(s) (ConfigMap)",
composed.config_mounts.len()
));
}
if !composed.loaded_env_files.is_empty() {
actions.push(format!(
"loaded {} env file(s)",
composed.loaded_env_files.len()
));
}
if !composed.missing_required.is_empty() {
actions.push(format!(
"WARNING: missing required env: {}",
composed.missing_required.join(", ")
));
}
if workload.is_some() {
actions.push("wait rollout".into());
}
if let Some(ex) = &env_cfg.expose {
let st = ex.service_type.as_deref().unwrap_or("ClusterIP");
let mut bits = vec![format!("expose type={st}")];
if let Some(lp) = ex.local_port {
bits.push(format!("local_port={lp}"));
}
if ex.port_forward.unwrap_or(ex.local_port.is_some() || !ex.dns_hosts.is_empty()) {
bits.push("port-forward".into());
}
if !ex.dns_hosts.is_empty() {
bits.push(format!("dns_hosts={}", ex.dns_hosts.len()));
}
actions.push(bits.join(", "));
} else if flags.local_image {
actions.push(
"hint: set deploy.envs.*.expose (local_port/port_forward) for localhost routing"
.into(),
);
}
if !env_cfg.health.is_empty() {
actions.push(format!("probe {} health URL(s)", env_cfg.health.len()));
}
}
} else if is_cloudflare_provider(&provider) {
let app = worker_app.as_deref().unwrap_or(svc.name.as_str());
actions.push(format!("cloudflare doctor/deploy worker `{app}`"));
if dockerfile.is_some() && image_ref.is_some() {
actions.push(
"skip GHCR build/push (Cloudflare owns image — OpenNext or containers local_build)"
.into(),
);
}
if let Some(r) = &rollout {
actions.push(format!("containers rollout `{r}`"));
} else {
actions.push("containers rollout (worker default)".into());
}
if !env_cfg.health.is_empty() {
actions.push(format!("probe {} health URL(s)", env_cfg.health.len()));
}
} else if crate::providers::is_local_provider(&provider) {
actions.push("local preflight steps".into());
} else if is_railway_provider(&provider) {
actions.push("railway deploy (not wired — use cloudflare or kubernetes for now)".into());
} else if crate::providers::is_oci_provider(&provider) {
actions.push("oci promote/publish".into());
} else {
actions.push(format!("provider `{provider}` (no runner yet)"));
}
let mut runtime_env = composed.env.clone();
if let Some(port) = composed.container_port {
runtime_env
.entry("PORT".into())
.or_insert_with(|| port.to_string());
}
Ok(ServicePlan {
name: svc.name.clone(),
provider,
destination: Some(destination),
root_directory: svc.root_directory.clone(),
version: version.to_string(),
image,
image_ref,
digest: None,
dockerfile,
build_context,
platforms,
worker_app,
rollout,
runtime_env,
container_port: composed.container_port,
config_mounts: composed.config_mounts.clone(),
expose: env_cfg.expose.clone(),
deploy: ServiceDeployPlan {
namespace,
workload,
service: k8s_service,
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)
}
#[cfg(test)]
mod multi_destination_tests {
use super::*;
use crate::context::DeployFlags;
use crate::types::{ServiceDeployDestinationView, ServiceDeployView};
use std::collections::HashMap;
fn dual_target_deploy() -> ServiceDeployView {
let mut destinations = HashMap::new();
destinations.insert(
"kubernetes".into(),
ServiceDeployDestinationView {
provider: Some("kubernetes".into()),
worker: None,
rollout: None,
enabled: None,
},
);
destinations.insert(
"cloudflare".into(),
ServiceDeployDestinationView {
provider: Some("cloudflare-containers".into()),
worker: Some("athena-auth".into()),
rollout: Some("immediate".into()),
enabled: None,
},
);
ServiceDeployView {
provider: "cloudflare-containers".into(),
worker: Some("athena-auth".into()),
rollout: Some("immediate".into()),
destinations,
envs: HashMap::new(),
}
}
#[test]
fn default_uses_all_enabled_destinations() {
let deploy = dual_target_deploy();
let flags = DeployFlags::default();
let specs = resolve_destination_specs(&deploy, "athena-auth", &flags).unwrap();
assert_eq!(specs.len(), 2);
assert_eq!(specs[0].destination, "kubernetes");
assert!(is_kubernetes_provider(&specs[0].provider));
assert_eq!(specs[1].destination, "cloudflare");
assert!(is_cloudflare_provider(&specs[1].provider));
assert_eq!(specs[1].worker.as_deref(), Some("athena-auth"));
assert_eq!(specs[1].rollout.as_deref(), Some("immediate"));
}
#[test]
fn to_cloudflare_narrows() {
let deploy = dual_target_deploy();
let mut flags = DeployFlags::default();
flags.destinations = vec!["cloudflare".into()];
let specs = resolve_destination_specs(&deploy, "athena-auth", &flags).unwrap();
assert_eq!(specs.len(), 1);
assert_eq!(specs[0].destination, "cloudflare");
}
#[test]
fn to_all_expands() {
let deploy = dual_target_deploy();
let mut flags = DeployFlags::default();
flags.destinations = vec![DESTINATION_ALL.into()];
let specs = resolve_destination_specs(&deploy, "athena-auth", &flags).unwrap();
assert_eq!(specs.len(), 2);
}
#[test]
fn empty_destinations_map_uses_default_provider() {
let deploy = ServiceDeployView {
provider: "kubernetes".into(),
worker: None,
rollout: None,
destinations: HashMap::new(),
envs: HashMap::new(),
};
let flags = DeployFlags::default();
let specs = resolve_destination_specs(&deploy, "athena", &flags).unwrap();
assert_eq!(specs.len(), 1);
assert_eq!(specs[0].destination, "kubernetes");
assert!(is_kubernetes_provider(&specs[0].provider));
}
#[test]
fn disabled_destination_skipped_unless_explicit() {
let mut deploy = dual_target_deploy();
deploy
.destinations
.get_mut("kubernetes")
.unwrap()
.enabled = Some(false);
let flags = DeployFlags::default();
let specs = resolve_destination_specs(&deploy, "athena-auth", &flags).unwrap();
assert_eq!(specs.len(), 1);
assert_eq!(specs[0].destination, "cloudflare");
let mut flags = DeployFlags::default();
flags.destinations = vec!["kubernetes".into()];
let specs = resolve_destination_specs(&deploy, "athena-auth", &flags).unwrap();
assert_eq!(specs.len(), 1);
assert_eq!(specs[0].destination, "kubernetes");
}
}