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 protocol_major: artifact.protocol_major,
254 entry: artifact.entry.clone(),
255 entries: artifact.entries.clone(),
256 style_assets: artifact.style_assets.clone(),
257 manifest: artifact.manifest.clone(),
258 requested_permissions: artifact.requested_permissions.clone(),
259 }
260 })
261 .collect();
262 effects.push(ModulePlanEffect::ConsoleComposition {
263 effect_id: "25-console-composition:lenso-console".to_owned(),
264 console_service_id: "lenso-console".to_owned(),
265 candidate_lock_digest: target_lock_digest.to_owned(),
266 artifacts,
267 });
268 }
269 let current_services = current_lock
270 .into_iter()
271 .flat_map(|lock| &lock.modules)
272 .filter_map(|module| match &module.delivery {
273 ModuleDelivery::Service(service) => Some((
274 service.service_id.clone(),
275 service.service_release_digest.clone(),
276 )),
277 ModuleDelivery::Linked(_) => None,
278 })
279 .collect::<BTreeSet<_>>();
280 let target_services = target_lock
281 .modules
282 .iter()
283 .filter_map(|module| match &module.delivery {
284 ModuleDelivery::Service(service) => Some((
285 service.service_id.clone(),
286 service.service_release_digest.clone(),
287 )),
288 ModuleDelivery::Linked(_) => None,
289 })
290 .collect::<BTreeSet<_>>();
291 let mut installation_state = current_service_installations.clone();
292 for (service_id, service_release_digest) in target_services.difference(¤t_services) {
293 let binding = service_binding(service_deployments, service_id, service_release_digest);
294 if let Some(installation) = binding.and_then(|binding| binding.installation.as_ref()) {
295 if installation.service_ref.service_id != *service_id
296 || installation.service_release.digest != *service_release_digest
297 {
298 return Err(crate::ServiceInstallationError::InvalidContract(
299 "Service Installation binding differs from the resolved Service release"
300 .to_owned(),
301 )
302 .into());
303 }
304 let installed_exports = installation
305 .exports
306 .iter()
307 .map(|export| export.module_id.as_str())
308 .collect::<BTreeSet<_>>();
309 if target_lock.modules.iter().any(|module| {
310 matches!(
311 &module.delivery,
312 ModuleDelivery::Service(service)
313 if service.service_id == *service_id
314 && service.service_release_digest == *service_release_digest
315 && !installed_exports.contains(module.module_id.as_str())
316 )
317 }) {
318 return Err(crate::ServiceInstallationError::InvalidContract(
319 "Service Installation does not declare every resolved Module export".to_owned(),
320 )
321 .into());
322 }
323 }
324 let installation_plan = binding
325 .and_then(|binding| binding.installation.clone())
326 .map(|installation| {
327 crate::plan_service_installation(
328 &installation_state,
329 crate::ServiceInstallationChange::Install { installation },
330 created_at,
331 )
332 })
333 .transpose()?;
334 if let Some(plan) = &installation_plan {
335 installation_state = plan.target.clone();
336 }
337 effects.push(ModulePlanEffect::ServiceInstallation {
338 effect_id: format!(
339 "20-service-install:{}:{}",
340 safe_id(service_id),
341 &service_release_digest[7..23]
342 ),
343 service_id: service_id.clone(),
344 service_release_digest: service_release_digest.clone(),
345 installation_plan,
346 adapter: binding.map(|binding| binding.adapter),
347 action: binding.and_then(|binding| binding.install.clone()),
348 });
349 }
350 for module in &changed_target_modules {
351 match &module.delivery {
352 ModuleDelivery::Service(_) => {}
353 ModuleDelivery::Linked(_) => {
354 if current_modules
355 .get(module.module_id.as_str())
356 .is_some_and(|current| current.release_digest == module.release_digest)
357 {
358 continue;
359 }
360 let Some(release) = releases.get(module.release_digest.as_str()) else {
361 continue;
362 };
363 for (declaration, artifact) in release
364 .manifest
365 .migrations
366 .iter()
367 .zip(&module.migration_artifacts)
368 {
369 effects.push(ModulePlanEffect::Migration {
370 effect_id: format!(
371 "30-migration:{}:{:08}:{}",
372 safe_id(&module.module_id),
373 declaration.order,
374 safe_id(&declaration.migration_id)
375 ),
376 module_id: module.module_id.clone(),
377 release_digest: module.release_digest.clone(),
378 migration_id: declaration.migration_id.clone(),
379 artifact_locator: artifact.locator.clone(),
380 artifact_digest: artifact.digest.clone(),
381 store_scope: declaration.store.clone(),
382 execution: MigrationExecutionMode::Transactional,
383 risk_class: if declaration.destructive {
384 ModuleRiskClass::DestructiveMigration
385 } else {
386 ModuleRiskClass::Ordinary
387 },
388 });
389 }
390 }
391 }
392 }
393 effects.push(ModulePlanEffect::Validate {
394 effect_id: "80-validate:cargo-check".to_owned(),
395 command: "cargo check --locked".to_owned(),
396 expected_evidence: cargo_lock_digest.to_owned(),
397 });
398 if changed_target_modules
399 .iter()
400 .any(|module| matches!(module.delivery, ModuleDelivery::Linked(_)))
401 {
402 effects.push(ModulePlanEffect::Restart {
403 effect_id: "99-restart:host".to_owned(),
404 target: "host".to_owned(),
405 });
406 }
407 for (service_id, service_release_digest) in &target_services {
408 if !changed_target_modules.iter().any(|module| {
409 matches!(&module.delivery, ModuleDelivery::Service(service) if service.service_id == *service_id && service.service_release_digest == *service_release_digest)
410 }) {
411 continue;
412 }
413 let binding = service_binding(service_deployments, service_id, service_release_digest);
414 let Some((adapter, action)) = binding.and_then(|binding| {
415 binding
416 .restart
417 .clone()
418 .map(|action| (binding.adapter, action))
419 }) else {
420 continue;
421 };
422 effects.push(ModulePlanEffect::ServiceRestart {
423 effect_id: format!("99-service-restart:{}", safe_id(service_id)),
424 service_id: service_id.clone(),
425 service_release_digest: service_release_digest.clone(),
426 adapter: Some(adapter),
427 action: Some(action),
428 });
429 }
430 effects.push(ModulePlanEffect::Activate {
431 effect_id: "90-activate:application-lock".to_owned(),
432 target_lock_digest: target_lock_digest.to_owned(),
433 });
434 effects.sort_by(|left, right| left.effect_id().cmp(right.effect_id()));
435 Ok(effects)
436}
437
438fn service_binding<'a>(
439 bindings: &'a [crate::ServiceDeploymentBinding],
440 service_id: &str,
441 release_digest: &str,
442) -> Option<&'a crate::ServiceDeploymentBinding> {
443 bindings.iter().find(|binding| {
444 binding.service_id == service_id && binding.service_release_digest == release_digest
445 })
446}
447
448fn destructive_boundaries(effects: &[ModulePlanEffect]) -> Vec<ModuleApprovalBoundary> {
449 let effect_ids = effects
450 .iter()
451 .filter(|effect| effect.risk_class() == ModuleRiskClass::DestructiveMigration)
452 .map(|effect| effect.effect_id().to_owned())
453 .collect::<Vec<_>>();
454 if effect_ids.is_empty() {
455 Vec::new()
456 } else {
457 vec![ModuleApprovalBoundary {
458 boundary_id: "destructive-migrations".to_owned(),
459 risk_class: ModuleRiskClass::DestructiveMigration,
460 required_authority: "module.migrate.destructive".to_owned(),
461 effect_ids,
462 backup_evidence_digest: None,
463 }]
464 }
465}
466
467fn safe_id(value: &str) -> String {
468 value
469 .chars()
470 .map(|character| {
471 if character.is_ascii_alphanumeric() || matches!(character, '-' | '_' | '.') {
472 character
473 } else {
474 '-'
475 }
476 })
477 .collect()
478}