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