1use crate::{
2 ApplicationModuleLock, CargoLockGenerator, DesiredModuleComposition, IsolatedCargoLockResolver,
3 LinkedWorkspaceError, LinkedWorkspacePlanner, MODULE_CHANGE_PLAN_PROTOCOL,
4 MigrationExecutionMode, ModuleApprovalBoundary, ModuleChangePlan, ModuleGraphResolver,
5 ModuleManagementError, ModulePlanEffect, ModuleResolutionCandidate, ModuleResolutionError,
6 ModuleResolutionRequest, ModuleRiskClass, ModuleRootChange, application_module_lock_digest,
7 desired_composition_digest, module_change_plan_digest, validate_change_plan,
8};
9use chrono::{DateTime, Utc};
10use lenso_contracts::{ModuleDelivery, digest_json};
11use serde_json::json;
12use std::collections::{BTreeMap, BTreeSet};
13use std::path::PathBuf;
14
15#[derive(Debug, Clone)]
16pub struct ModuleChangePlanRequest {
17 pub current_desired: DesiredModuleComposition,
18 pub current_lock: Option<ApplicationModuleLock>,
19 pub change: ModuleRootChange,
20 pub catalog_snapshot_digest: String,
21 pub trust_policy_digest: String,
22 pub compatibility_evidence_digest: String,
23 pub resolver_version: String,
24 pub environment_id: String,
25 pub expected_target_revision: u64,
26 pub candidates: Vec<ModuleResolutionCandidate>,
27 pub current_service_installations: crate::ServiceInstallationSet,
28 pub service_deployments: Vec<crate::ServiceDeploymentBinding>,
29 pub cargo_offline: bool,
30 pub created_at: DateTime<Utc>,
31}
32
33#[derive(Debug, thiserror::Error)]
34pub enum ModuleChangePlannerError {
35 #[error(transparent)]
36 Resolution(#[from] ModuleResolutionError),
37 #[error(transparent)]
38 Linked(#[from] LinkedWorkspaceError),
39 #[error(transparent)]
40 Cargo(#[from] crate::CargoLockResolutionError),
41 #[error(transparent)]
42 Management(#[from] ModuleManagementError),
43 #[error("Module change plan JSON failed: {0}")]
44 Json(#[from] serde_json::Error),
45 #[error(transparent)]
46 ServiceInstallation(#[from] crate::ServiceInstallationError),
47 #[error("isolated Cargo produced a non-UTF-8 lockfile")]
48 NonUtf8CargoLock,
49}
50
51#[derive(Debug, Clone)]
52pub struct ModuleChangePlanner<G = crate::CargoGenerateLockfile> {
53 graph: ModuleGraphResolver,
54 linked: LinkedWorkspacePlanner,
55 cargo: IsolatedCargoLockResolver<G>,
56}
57
58impl ModuleChangePlanner<crate::CargoGenerateLockfile> {
59 pub fn new(workspace_root: impl Into<PathBuf>) -> Self {
60 Self {
61 graph: ModuleGraphResolver,
62 linked: LinkedWorkspacePlanner::new(workspace_root),
63 cargo: IsolatedCargoLockResolver::default(),
64 }
65 }
66}
67
68impl<G> ModuleChangePlanner<G>
69where
70 G: CargoLockGenerator,
71{
72 pub fn with_cargo_generator(workspace_root: impl Into<PathBuf>, generator: G) -> Self {
73 Self {
74 graph: ModuleGraphResolver,
75 linked: LinkedWorkspacePlanner::new(workspace_root),
76 cargo: IsolatedCargoLockResolver::new(generator),
77 }
78 }
79
80 pub fn plan(
81 &self,
82 request: &ModuleChangePlanRequest,
83 ) -> Result<ModuleChangePlan, ModuleChangePlannerError> {
84 let resolution = self.graph.resolve(&ModuleResolutionRequest {
85 current_desired: request.current_desired.clone(),
86 current_lock: request.current_lock.clone(),
87 change: request.change.clone(),
88 catalog_snapshot_digest: request.catalog_snapshot_digest.clone(),
89 trust_policy_digest: request.trust_policy_digest.clone(),
90 resolver_version: request.resolver_version.clone(),
91 candidates: request.candidates.clone(),
92 })?;
93 let reviewed_desired_document = format!(
94 "{}\n",
95 serde_json::to_string_pretty(&resolution.target_desired)?
96 );
97 let cargo_preparation = self.linked.prepare_cargo_resolution(
98 &resolution.target_desired,
99 request.current_lock.as_ref(),
100 &resolution.target_lock,
101 &reviewed_desired_document,
102 request.cargo_offline,
103 )?;
104 let cargo = self.cargo.resolve(&cargo_preparation.cargo_request)?;
105 let candidate_cargo_lock = String::from_utf8(cargo.candidate_lock)
106 .map_err(|_| ModuleChangePlannerError::NonUtf8CargoLock)?;
107 let workspace = self.linked.plan(
108 &resolution.target_desired,
109 &resolution.target_lock,
110 &reviewed_desired_document,
111 Some(&candidate_cargo_lock),
112 )?;
113
114 let current_desired_digest = desired_composition_digest(&request.current_desired)?;
115 let target_desired_digest = desired_composition_digest(&resolution.target_desired)?;
116 let current_lock_digest = request
117 .current_lock
118 .as_ref()
119 .map(application_module_lock_digest)
120 .transpose()?;
121 let target_lock_digest = application_module_lock_digest(&resolution.target_lock)?;
122 let plan_identity = digest_json(&json!({
123 "application_id": resolution.target_desired.application_id,
124 "environment_id": request.environment_id,
125 "expected_target_revision": request.expected_target_revision,
126 "request": request.change,
127 "current_desired_digest": current_desired_digest,
128 "current_lock_digest": current_lock_digest,
129 "target_lock_digest": target_lock_digest,
130 }))?;
131 let mut effects = workspace.effects;
132 effects.extend(non_workspace_effects(
133 request.current_lock.as_ref(),
134 &resolution.target_lock,
135 &request.candidates,
136 &target_lock_digest,
137 &cargo.evidence.candidate_lock_digest,
138 &request.current_service_installations,
139 &request.service_deployments,
140 request.created_at,
141 )?);
142 effects.sort_by(|left, right| left.effect_id().cmp(right.effect_id()));
143 let approval_boundaries = destructive_boundaries(&effects);
144 let mut next_actions = vec!["review_plan".to_owned(), "apply_plan".to_owned()];
145 if !approval_boundaries.is_empty() {
146 next_actions.insert(1, "approve_destructive_migrations".to_owned());
147 }
148 let validation_commands = vec!["cargo check --locked".to_owned()];
149 let mut plan = ModuleChangePlan {
150 protocol: MODULE_CHANGE_PLAN_PROTOCOL.to_owned(),
151 plan_id: format!("module-plan-{}", &plan_identity[7..23]),
152 plan_digest: String::new(),
153 application_id: resolution.target_desired.application_id.clone(),
154 environment_id: request.environment_id.clone(),
155 expected_target_revision: request.expected_target_revision,
156 request: request.change.clone(),
157 current_desired_digest,
158 target_desired: resolution.target_desired,
159 target_desired_digest,
160 current_lock_digest,
161 target_lock: resolution.target_lock,
162 target_lock_digest,
163 catalog_snapshot_digest: request.catalog_snapshot_digest.clone(),
164 resolver_version: request.resolver_version.clone(),
165 trust_policy_digest: request.trust_policy_digest.clone(),
166 compatibility_evidence_digest: request.compatibility_evidence_digest.clone(),
167 cargo_lock_candidate: Some(cargo.evidence),
168 read_set: workspace.read_set,
169 effects,
170 approval_boundaries,
171 validation_commands,
172 next_actions,
173 created_at: request.created_at,
174 };
175 plan.plan_digest = module_change_plan_digest(&plan)?;
176 validate_change_plan(&plan)?;
177 Ok(plan)
178 }
179}
180
181#[allow(clippy::too_many_lines)]
182fn non_workspace_effects(
183 current_lock: Option<&ApplicationModuleLock>,
184 target_lock: &ApplicationModuleLock,
185 candidates: &[ModuleResolutionCandidate],
186 target_lock_digest: &str,
187 cargo_lock_digest: &str,
188 current_service_installations: &crate::ServiceInstallationSet,
189 service_deployments: &[crate::ServiceDeploymentBinding],
190 created_at: DateTime<Utc>,
191) -> Result<Vec<ModulePlanEffect>, ModuleChangePlannerError> {
192 let releases = candidates
193 .iter()
194 .map(|candidate| (candidate.release_digest.as_str(), &candidate.release))
195 .collect::<BTreeMap<_, _>>();
196 let mut effects = Vec::new();
197 let current_modules = current_lock
198 .map(|lock| {
199 lock.modules
200 .iter()
201 .map(|module| (module.module_id.as_str(), module))
202 .collect::<BTreeMap<_, _>>()
203 })
204 .unwrap_or_default();
205 let changed_target_modules = target_lock
206 .modules
207 .iter()
208 .filter(|module| {
209 current_modules
210 .get(module.module_id.as_str())
211 .is_none_or(|current| {
212 current.release_digest != module.release_digest
213 || current.delivery != module.delivery
214 || current.crate_features != module.crate_features
215 || current.local_override_digest != module.local_override_digest
216 })
217 })
218 .collect::<Vec<_>>();
219 let current_console_artifacts = current_lock
220 .into_iter()
221 .flat_map(|lock| &lock.modules)
222 .filter_map(|module| {
223 module.console_ui_artifact.as_ref().map(|artifact| {
224 (
225 module.module_id.as_str(),
226 (module.release_digest.as_str(), artifact),
227 )
228 })
229 })
230 .collect::<BTreeMap<_, _>>();
231 let target_console_artifacts = target_lock
232 .modules
233 .iter()
234 .filter_map(|module| {
235 module.console_ui_artifact.as_ref().map(|artifact| {
236 (
237 module.module_id.as_str(),
238 (module.release_digest.as_str(), artifact),
239 )
240 })
241 })
242 .collect::<BTreeMap<_, _>>();
243 if current_console_artifacts != target_console_artifacts {
244 let artifacts = target_console_artifacts
245 .into_iter()
246 .map(|(module_id, (module_release_digest, artifact))| {
247 crate::ConsoleCompositionArtifact {
248 module_id: module_id.to_owned(),
249 module_release_digest: module_release_digest.to_owned(),
250 locator: artifact.locator.clone(),
251 digest: artifact.digest.clone(),
252 format: artifact.format.clone(),
253 entries: artifact.entries.clone(),
254 bridge_protocol: artifact.bridge_protocol.clone(),
255 requested_permissions: artifact.requested_permissions.clone(),
256 }
257 })
258 .collect();
259 effects.push(ModulePlanEffect::ConsoleComposition {
260 effect_id: "25-console-composition:lenso-console".to_owned(),
261 console_service_id: "lenso-console".to_owned(),
262 candidate_lock_digest: target_lock_digest.to_owned(),
263 artifacts,
264 });
265 }
266 let current_services = current_lock
267 .into_iter()
268 .flat_map(|lock| &lock.modules)
269 .filter_map(|module| match &module.delivery {
270 ModuleDelivery::Service(service) => Some((
271 service.service_id.clone(),
272 service.service_release_digest.clone(),
273 )),
274 ModuleDelivery::Linked(_) => None,
275 })
276 .collect::<BTreeSet<_>>();
277 let target_services = target_lock
278 .modules
279 .iter()
280 .filter_map(|module| match &module.delivery {
281 ModuleDelivery::Service(service) => Some((
282 service.service_id.clone(),
283 service.service_release_digest.clone(),
284 )),
285 ModuleDelivery::Linked(_) => None,
286 })
287 .collect::<BTreeSet<_>>();
288 let mut installation_state = current_service_installations.clone();
289 for (service_id, service_release_digest) in target_services.difference(¤t_services) {
290 let binding = service_binding(service_deployments, service_id, service_release_digest);
291 if let Some(installation) = binding.and_then(|binding| binding.installation.as_ref()) {
292 if installation.service_ref.service_id != *service_id
293 || installation.service_release.digest != *service_release_digest
294 {
295 return Err(crate::ServiceInstallationError::InvalidContract(
296 "Service Installation binding differs from the resolved Service release"
297 .to_owned(),
298 )
299 .into());
300 }
301 let installed_exports = installation
302 .exports
303 .iter()
304 .map(|export| export.module_id.as_str())
305 .collect::<BTreeSet<_>>();
306 if target_lock.modules.iter().any(|module| {
307 matches!(
308 &module.delivery,
309 ModuleDelivery::Service(service)
310 if service.service_id == *service_id
311 && service.service_release_digest == *service_release_digest
312 && !installed_exports.contains(module.module_id.as_str())
313 )
314 }) {
315 return Err(crate::ServiceInstallationError::InvalidContract(
316 "Service Installation does not declare every resolved Module export".to_owned(),
317 )
318 .into());
319 }
320 }
321 let installation_plan = binding
322 .and_then(|binding| binding.installation.clone())
323 .map(|installation| {
324 crate::plan_service_installation(
325 &installation_state,
326 crate::ServiceInstallationChange::Install { installation },
327 created_at,
328 )
329 })
330 .transpose()?;
331 if let Some(plan) = &installation_plan {
332 installation_state = plan.target.clone();
333 }
334 effects.push(ModulePlanEffect::ServiceInstallation {
335 effect_id: format!(
336 "20-service-install:{}:{}",
337 safe_id(service_id),
338 &service_release_digest[7..23]
339 ),
340 service_id: service_id.clone(),
341 service_release_digest: service_release_digest.clone(),
342 installation_plan,
343 adapter: binding.map(|binding| binding.adapter),
344 action: binding.and_then(|binding| binding.install.clone()),
345 });
346 }
347 for module in &changed_target_modules {
348 match &module.delivery {
349 ModuleDelivery::Service(_) => {}
350 ModuleDelivery::Linked(_) => {
351 if current_modules
352 .get(module.module_id.as_str())
353 .is_some_and(|current| current.release_digest == module.release_digest)
354 {
355 continue;
356 }
357 let Some(release) = releases.get(module.release_digest.as_str()) else {
358 continue;
359 };
360 for (declaration, artifact) in release
361 .manifest
362 .migrations
363 .iter()
364 .zip(&module.migration_artifacts)
365 {
366 effects.push(ModulePlanEffect::Migration {
367 effect_id: format!(
368 "30-migration:{}:{:08}:{}",
369 safe_id(&module.module_id),
370 declaration.order,
371 safe_id(&declaration.migration_id)
372 ),
373 module_id: module.module_id.clone(),
374 release_digest: module.release_digest.clone(),
375 migration_id: declaration.migration_id.clone(),
376 artifact_locator: artifact.locator.clone(),
377 artifact_digest: artifact.digest.clone(),
378 store_scope: declaration.store.clone(),
379 execution: MigrationExecutionMode::Transactional,
380 risk_class: if declaration.destructive {
381 ModuleRiskClass::DestructiveMigration
382 } else {
383 ModuleRiskClass::Ordinary
384 },
385 });
386 }
387 }
388 }
389 }
390 effects.push(ModulePlanEffect::Validate {
391 effect_id: "80-validate:cargo-check".to_owned(),
392 command: "cargo check --locked".to_owned(),
393 expected_evidence: cargo_lock_digest.to_owned(),
394 });
395 if changed_target_modules
396 .iter()
397 .any(|module| matches!(module.delivery, ModuleDelivery::Linked(_)))
398 {
399 effects.push(ModulePlanEffect::Restart {
400 effect_id: "99-restart:host".to_owned(),
401 target: "host".to_owned(),
402 });
403 }
404 for (service_id, service_release_digest) in &target_services {
405 if !changed_target_modules.iter().any(|module| {
406 matches!(&module.delivery, ModuleDelivery::Service(service) if service.service_id == *service_id && service.service_release_digest == *service_release_digest)
407 }) {
408 continue;
409 }
410 let binding = service_binding(service_deployments, service_id, service_release_digest);
411 let Some((adapter, action)) = binding.and_then(|binding| {
412 binding
413 .restart
414 .clone()
415 .map(|action| (binding.adapter, action))
416 }) else {
417 continue;
418 };
419 effects.push(ModulePlanEffect::ServiceRestart {
420 effect_id: format!("99-service-restart:{}", safe_id(service_id)),
421 service_id: service_id.clone(),
422 service_release_digest: service_release_digest.clone(),
423 adapter: Some(adapter),
424 action: Some(action),
425 });
426 }
427 effects.push(ModulePlanEffect::Activate {
428 effect_id: "90-activate:application-lock".to_owned(),
429 target_lock_digest: target_lock_digest.to_owned(),
430 });
431 effects.sort_by(|left, right| left.effect_id().cmp(right.effect_id()));
432 Ok(effects)
433}
434
435fn service_binding<'a>(
436 bindings: &'a [crate::ServiceDeploymentBinding],
437 service_id: &str,
438 release_digest: &str,
439) -> Option<&'a crate::ServiceDeploymentBinding> {
440 bindings.iter().find(|binding| {
441 binding.service_id == service_id && binding.service_release_digest == release_digest
442 })
443}
444
445fn destructive_boundaries(effects: &[ModulePlanEffect]) -> Vec<ModuleApprovalBoundary> {
446 let effect_ids = effects
447 .iter()
448 .filter(|effect| effect.risk_class() == ModuleRiskClass::DestructiveMigration)
449 .map(|effect| effect.effect_id().to_owned())
450 .collect::<Vec<_>>();
451 if effect_ids.is_empty() {
452 Vec::new()
453 } else {
454 vec![ModuleApprovalBoundary {
455 boundary_id: "destructive-migrations".to_owned(),
456 risk_class: ModuleRiskClass::DestructiveMigration,
457 required_authority: "module.migrate.destructive".to_owned(),
458 effect_ids,
459 backup_evidence_digest: None,
460 }]
461 }
462}
463
464fn safe_id(value: &str) -> String {
465 value
466 .chars()
467 .map(|character| {
468 if character.is_ascii_alphanumeric() || matches!(character, '-' | '_' | '.') {
469 character
470 } else {
471 '-'
472 }
473 })
474 .collect()
475}