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 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 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 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}