Skip to main content

lenso_module_management/
planner.rs

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(&current_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}