use crate::{
AdvanceModuleOperation, ApplicationModuleLock, CargoLockResolutionRequest,
DesiredModuleComposition, ExpectedLinkedPackage, LINKED_COMPOSITION_SEAM_PROTOCOL,
LinkedCompositionSeam, ManagedFileType, ModuleEffectOutcome, ModuleEffectReceipt,
ModuleFileChange, ModuleFileOwnership, ModuleManagementEngine, ModuleManagementError,
ModuleOperation, ModuleOperationError, ModuleOperationState, ModuleOperationStore,
ModuleOperationStoreError, ModulePathPrecondition, ModulePlanEffect, ModuleWorkspaceBackup,
PathExistence, application_module_lock_digest, validate_application_module_lock,
validate_change_plan, validate_desired_composition,
};
use base64::{Engine as _, engine::general_purpose::STANDARD as BASE64};
use chrono::{DateTime, Utc};
use lenso_contracts::{ArtifactReference, ModuleDelivery, digest_json};
use sha2::{Digest as _, Sha256};
use std::collections::{BTreeMap, BTreeSet};
use std::fs::{self, File, OpenOptions};
use std::io::Write as _;
use std::path::{Component, Path, PathBuf};
use thiserror::Error;
const GENERATED_MARKER: &str = "generated by lenso-module-management";
#[derive(Debug, Error)]
pub enum LinkedWorkspaceError {
#[error("Linked composition contract is invalid: {0}")]
InvalidContract(String),
#[error("workspace plan is stale at `{0}`")]
Stale(String),
#[error("workspace path is unsafe: `{0}`")]
UnsafePath(String),
#[error("workspace I/O failed: {0}")]
Io(#[from] std::io::Error),
#[error("workspace JSON failed: {0}")]
Json(#[from] serde_json::Error),
#[error(transparent)]
Management(#[from] ModuleManagementError),
#[error(transparent)]
Store(#[from] ModuleOperationStoreError),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct LinkedWorkspacePlan {
pub read_set: Vec<ModulePathPrecondition>,
pub effects: Vec<ModulePlanEffect>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct LinkedCargoPreparation {
pub workspace_plan_without_lock: LinkedWorkspacePlan,
pub cargo_request: CargoLockResolutionRequest,
}
#[derive(Debug, Clone)]
pub struct LinkedWorkspacePlanner {
root: PathBuf,
seam_path: PathBuf,
}
impl LinkedWorkspacePlanner {
pub fn new(root: impl Into<PathBuf>) -> Self {
Self {
root: root.into(),
seam_path: PathBuf::from(".lenso/linked-composition-seam.json"),
}
}
#[must_use]
pub fn with_seam_path(mut self, path: impl Into<PathBuf>) -> Self {
self.seam_path = path.into();
self
}
pub fn plan(
&self,
desired: &DesiredModuleComposition,
module_lock: &ApplicationModuleLock,
reviewed_desired_document: &str,
candidate_cargo_lock: Option<&str>,
) -> Result<LinkedWorkspacePlan, LinkedWorkspaceError> {
validate_desired_composition(desired)?;
validate_application_module_lock(module_lock)?;
let reviewed: DesiredModuleComposition = serde_json::from_str(reviewed_desired_document)?;
if &reviewed != desired {
return invalid("reviewed Desired Composition bytes do not match the target contract");
}
let seam_relative = normalized_relative(&self.seam_path)?;
let seam_snapshot = snapshot_path(&self.root, &seam_relative)?;
let seam_bytes = seam_snapshot.bytes.as_deref().ok_or_else(|| {
LinkedWorkspaceError::InvalidContract("stable Host seam is missing".into())
})?;
let seam: LinkedCompositionSeam = serde_json::from_slice(seam_bytes)?;
validate_seam(&self.root, &seam)?;
let lock_digest = application_module_lock_digest(module_lock)?;
let generated = generate_composition_files(desired, module_lock, &seam, &lock_digest)?;
let mut targets = vec![
PlannedFile::user("lenso.modules.json", reviewed_desired_document.as_bytes()),
PlannedFile::generated_json(
"lenso.modules.lock.json",
format!("{}\n", serde_json::to_string_pretty(module_lock)?).as_bytes(),
),
];
targets.extend(generated);
if let Some(candidate) = candidate_cargo_lock {
if candidate.trim().is_empty() {
return invalid("candidate Cargo.lock must not be empty");
}
targets.push(PlannedFile::user("Cargo.lock", candidate.as_bytes()));
}
targets.sort_by(|left, right| left.path.cmp(&right.path));
if !targets.windows(2).all(|pair| pair[0].path < pair[1].path) {
return invalid("Linked composition target paths collide");
}
let mut read_paths = BTreeSet::from([
seam_relative,
seam.host_manifest_path.clone(),
seam.host_source_path.clone(),
]);
read_paths.extend(targets.iter().map(|target| target.path.clone()));
let snapshots = read_paths
.iter()
.map(|path| snapshot_path(&self.root, path).map(|snapshot| (path.clone(), snapshot)))
.collect::<Result<BTreeMap<_, _>, _>>()?;
let read_set = snapshots
.values()
.map(|snapshot| snapshot.precondition.clone())
.collect::<Vec<_>>();
let effects = targets
.into_iter()
.filter_map(|target| {
let before = snapshots
.get(&target.path)
.expect("target path was added to read set");
let after_digest = raw_digest(&target.content);
if before.precondition.content_digest.as_deref() == Some(after_digest.as_str()) {
return None;
}
let change = if before.precondition.existence == PathExistence::Absent {
ModuleFileChange::Create
} else {
ModuleFileChange::Modify
};
let before_display = before.bytes.as_deref().map(String::from_utf8_lossy);
Some(ModulePlanEffect::WorkspaceFile {
effect_id: format!("10-workspace:{}", target.path),
path: target.path.clone(),
ownership: target.ownership,
change,
before_digest: before.precondition.content_digest.clone(),
after_digest: Some(after_digest),
after_content: Some(
String::from_utf8(target.content)
.expect("managed Linked composition files are UTF-8"),
),
after_mode: Some(0o644),
patch: exact_replacement_patch(
&target.path,
before_display.as_deref(),
&target.display,
),
reversible_before_migration: true,
})
})
.collect::<Vec<_>>();
Ok(LinkedWorkspacePlan { read_set, effects })
}
pub fn prepare_cargo_resolution(
&self,
desired: &DesiredModuleComposition,
current_lock: Option<&ApplicationModuleLock>,
target_lock: &ApplicationModuleLock,
reviewed_desired_document: &str,
offline: bool,
) -> Result<LinkedCargoPreparation, LinkedWorkspaceError> {
let workspace_plan_without_lock =
self.plan(desired, target_lock, reviewed_desired_document, None)?;
let seam_relative = normalized_relative(&self.seam_path)?;
let seam: LinkedCompositionSeam = serde_json::from_slice(
snapshot_path(&self.root, &seam_relative)?
.bytes
.as_deref()
.ok_or_else(|| {
LinkedWorkspaceError::InvalidContract("stable Host seam is missing".into())
})?,
)?;
let root_manifest_path = seam.host_manifest_path.clone();
let manifest_parent = Path::new(&root_manifest_path)
.parent()
.filter(|path| !path.as_os_str().is_empty())
.unwrap_or_else(|| Path::new(""));
let lock_path = normalized_relative(&manifest_parent.join("Cargo.lock"))?;
let mut read_set = BTreeMap::new();
for path in [&root_manifest_path, &lock_path] {
let snapshot = snapshot_path(&self.root, path)?;
let bytes = snapshot.bytes.ok_or_else(|| {
LinkedWorkspaceError::InvalidContract(format!(
"isolated Cargo input `{path}` is missing"
))
})?;
read_set.insert(path.clone(), bytes);
}
let generated_prefix = format!("{}/", seam.generated_crate_path.trim_end_matches('/'));
let candidate_files = workspace_plan_without_lock
.effects
.iter()
.filter_map(|effect| match effect {
ModulePlanEffect::WorkspaceFile {
path,
after_content: Some(content),
..
} if path.starts_with(&generated_prefix) => {
Some((path.clone(), content.as_bytes().to_vec()))
}
_ => None,
})
.collect::<BTreeMap<_, _>>();
let current_linked_packages = current_lock
.map(expected_linked_packages)
.unwrap_or_default();
let expected_linked_packages = expected_linked_packages(target_lock);
let mut allowed_root_packages = current_linked_packages
.iter()
.chain(&expected_linked_packages)
.map(|package| package.package.clone())
.collect::<BTreeSet<_>>();
allowed_root_packages.insert("lenso-linked-composition".to_owned());
Ok(LinkedCargoPreparation {
workspace_plan_without_lock,
cargo_request: CargoLockResolutionRequest {
read_set,
candidate_files,
root_manifest_path,
lock_path,
allowed_root_packages: allowed_root_packages.into_iter().collect(),
current_linked_packages,
expected_linked_packages,
offline,
},
})
}
}
fn expected_linked_packages(module_lock: &ApplicationModuleLock) -> Vec<ExpectedLinkedPackage> {
let mut packages = module_lock
.modules
.iter()
.filter_map(|module| {
let ModuleDelivery::Linked(delivery) = &module.delivery else {
return None;
};
Some(ExpectedLinkedPackage {
package: delivery.package.clone(),
version: delivery.crate_version.clone(),
archive_checksum: module
.local_override_digest
.is_none()
.then(|| delivery.archive_checksum.clone()),
default_features: delivery.default_features,
features: module.crate_features.clone(),
})
})
.collect::<Vec<_>>();
packages.sort_by(|left, right| left.package.cmp(&right.package));
packages
}
#[derive(Debug, Clone)]
pub struct LinkedWorkspaceTransaction {
root: PathBuf,
}
impl LinkedWorkspaceTransaction {
pub fn new(root: impl Into<PathBuf>) -> Self {
Self { root: root.into() }
}
#[allow(clippy::too_many_arguments)]
pub fn apply<S: ModuleOperationStore>(
&self,
engine: &ModuleManagementEngine<S>,
operation_id: &str,
plan: &crate::ModuleChangePlan,
fencing_token: u64,
actor_id: &str,
now: DateTime<Utc>,
) -> Result<ModuleOperation, LinkedWorkspaceError> {
validate_change_plan(plan)?;
let mut operation = engine.store().load(operation_id)?;
if operation.plan_digest != plan.plan_digest {
return invalid("operation does not bind the supplied Linked workspace plan");
}
if operation.fencing_token != fencing_token {
return Err(LinkedWorkspaceError::Management(
ModuleManagementError::StaleFencingToken,
));
}
if operation.state == ModuleOperationState::FilesApplied {
return Ok(operation);
}
let effects = workspace_effects(plan);
if effects.is_empty() {
return invalid("Linked workspace plan contains no file effects");
}
if operation.state == ModuleOperationState::Ready {
verify_read_set(&self.root, plan, &effects, false)?;
let backups = collect_backups(&self.root, &effects)?;
operation = engine.begin_workspace_application(
operation_id,
operation.revision,
fencing_token,
actor_id,
plan,
backups,
Vec::new(),
now,
)?;
} else if operation.state == ModuleOperationState::ApplyingFiles {
if operation.workspace_backups.is_empty() {
return invalid("crash continuation has no journaled workspace backups");
}
verify_read_set(&self.root, plan, &effects, true)?;
} else {
return Err(LinkedWorkspaceError::Management(
ModuleManagementError::IllegalTransition {
from: operation.state,
to: ModuleOperationState::ApplyingFiles,
},
));
}
let result = self.apply_effects(engine, operation, effects, actor_id, now);
match result {
Ok(operation) => Ok(operation),
Err(error) => {
let current = engine.store().load(operation_id)?;
let restore = restore_backups(&self.root, ¤t.workspace_backups);
let (state, code, message) = match restore {
Ok(()) => (
ModuleOperationState::Restored,
"workspace_apply_failed_restored",
error.to_string(),
),
Err(restore_error) => (
ModuleOperationState::RepairRequired,
"workspace_restore_failed",
format!("{error}; restore failed: {restore_error}"),
),
};
engine.advance(AdvanceModuleOperation {
operation_id: current.operation_id,
expected_revision: current.revision,
fencing_token,
next_state: state,
actor_id: actor_id.to_owned(),
outcome_code: code.to_owned(),
evidence_references: Vec::new(),
error: Some(ModuleOperationError {
code: code.to_owned(),
message,
evidence_references: Vec::new(),
recorded_at: now,
}),
next_actions: vec![if state == ModuleOperationState::Restored {
"replan_from_restored_workspace".to_owned()
} else {
"create_workspace_repair_plan".to_owned()
}],
now,
})?;
Err(error)
}
}
}
pub fn resume_evidence(
&self,
operation: &ModuleOperation,
plan: &crate::ModuleChangePlan,
observed_at: DateTime<Utc>,
) -> Result<crate::ModuleResumeEvidence, LinkedWorkspaceError> {
validate_change_plan(plan)?;
if operation.state != ModuleOperationState::ApplyingFiles
|| operation.plan_digest != plan.plan_digest
{
return invalid("only a bound in-progress workspace transaction can resume");
}
let completed_effect_ids = operation
.effect_receipts
.iter()
.map(|receipt| receipt.effect_id.clone())
.collect::<Vec<_>>();
let completed = completed_effect_ids
.iter()
.map(String::as_str)
.collect::<BTreeSet<_>>();
let next = workspace_effects(plan)
.into_iter()
.find(|effect| !completed.contains(effect.effect_id()))
.ok_or_else(|| {
LinkedWorkspaceError::InvalidContract(
"workspace transaction has no incomplete effect".to_owned(),
)
})?;
let (path, _, _, _, _) = workspace_effect_parts(next)?;
let snapshot = snapshot_path(&self.root, path)?;
Ok(crate::ModuleResumeEvidence {
plan_digest: plan.plan_digest.clone(),
observed_target_digest: raw_digest(snapshot.bytes.as_deref().unwrap_or_default()),
completed_effect_ids,
next_effect_id: next.effect_id().to_owned(),
next_effect_idempotent: true,
observed_at,
})
}
fn apply_effects<S: ModuleOperationStore>(
&self,
engine: &ModuleManagementEngine<S>,
mut operation: ModuleOperation,
effects: Vec<&ModulePlanEffect>,
actor_id: &str,
now: DateTime<Utc>,
) -> Result<ModuleOperation, LinkedWorkspaceError> {
for effect in effects {
let effect_id = effect.effect_id();
if operation
.effect_receipts
.iter()
.any(|receipt| receipt.effect_id == effect_id)
{
continue;
}
let (path, change, after_digest, after_content, after_mode) =
workspace_effect_parts(effect)?;
let snapshot = snapshot_path(&self.root, path)?;
let outcome = if matches_after(&snapshot.precondition, change, after_digest) {
ModuleEffectOutcome::AlreadyApplied
} else {
apply_file(
&self.root,
path,
change,
after_content,
after_mode,
&operation,
effect_id,
)?;
ModuleEffectOutcome::Applied
};
let effect_digest = digest_json(effect)?;
operation = engine.record_effect_receipt(
&operation.operation_id,
operation.revision,
operation.fencing_token,
actor_id,
ModuleEffectReceipt {
receipt_id: format!("{}:{}", operation.operation_id, effect_id),
effect_id: effect_id.to_owned(),
effect_digest,
operation_id: operation.operation_id.clone(),
attempt: operation.attempt,
fencing_token: operation.fencing_token,
outcome,
evidence_references: after_digest
.map(|digest| ArtifactReference {
locator: path.to_owned(),
digest: digest.to_owned(),
})
.into_iter()
.collect(),
committed_at: now,
},
now,
)?;
}
Ok(engine.advance(AdvanceModuleOperation {
operation_id: operation.operation_id,
expected_revision: operation.revision,
fencing_token: operation.fencing_token,
next_state: ModuleOperationState::FilesApplied,
actor_id: actor_id.to_owned(),
outcome_code: "linked_workspace_files_applied".to_owned(),
evidence_references: Vec::new(),
error: None,
next_actions: vec!["run_reviewed_locked_validation".to_owned()],
now,
})?)
}
}
#[derive(Debug)]
struct PlannedFile {
path: String,
ownership: ModuleFileOwnership,
content: Vec<u8>,
display: String,
}
impl PlannedFile {
fn user(path: &str, content: &[u8]) -> Self {
Self::new(path, ModuleFileOwnership::User, content)
}
fn generated_json(path: &str, content: &[u8]) -> Self {
Self::new(path, ModuleFileOwnership::Generated, content)
}
fn generated(path: &str, content: &str) -> Self {
Self::new(path, ModuleFileOwnership::Generated, content.as_bytes())
}
fn new(path: &str, ownership: ModuleFileOwnership, content: &[u8]) -> Self {
Self {
path: path.to_owned(),
ownership,
content: content.to_vec(),
display: String::from_utf8_lossy(content).into_owned(),
}
}
}
#[derive(Debug)]
struct PathSnapshot {
precondition: ModulePathPrecondition,
bytes: Option<Vec<u8>>,
}
fn generate_composition_files(
desired: &DesiredModuleComposition,
module_lock: &ApplicationModuleLock,
seam: &LinkedCompositionSeam,
lock_digest: &str,
) -> Result<Vec<PlannedFile>, LinkedWorkspaceError> {
let overrides = desired
.local_overrides
.iter()
.map(|entry| (entry.module_id.as_str(), entry.path.as_str()))
.collect::<BTreeMap<_, _>>();
let mut aliases = BTreeSet::new();
let mut dependencies = Vec::new();
let modules = topological_linked_modules(module_lock)?;
let mut bindings = Vec::new();
for module in modules {
let ModuleDelivery::Linked(delivery) = &module.delivery else {
continue;
};
let alias = dependency_alias(&module.module_id);
if !aliases.insert(alias.clone()) {
return invalid("deterministic Linked dependency aliases collide");
}
if !valid_package_name(&delivery.package) || !valid_binding_path(&delivery.binding) {
return invalid("Linked package or binding export is not safe to generate");
}
if !module
.crate_features
.windows(2)
.all(|pair| pair[0] < pair[1])
{
return invalid("Linked crate features must be sorted and unique");
}
let features = module
.crate_features
.iter()
.map(|feature| format!("\"{}\"", escape_toml(feature)))
.collect::<Vec<_>>()
.join(", ");
let coordinate = overrides.get(module.module_id.as_str()).map_or_else(
|| format!("version = \"={}\"", delivery.crate_version),
|path| format!("path = \"{}\"", escape_toml(path)),
);
dependencies.push(format!(
"{alias} = {{ package = \"{}\", {coordinate}, default-features = {}, features = [{features}] }}",
escape_toml(&delivery.package),
delivery.default_features,
));
bindings.push(format!(" {alias}::{}(),", delivery.binding));
}
let cargo = format!(
"# {GENERATED_MARKER}; source-lock: {lock_digest}\n[package]\nname = \"lenso-linked-composition\"\nversion = \"0.0.0\"\nedition = \"2024\"\npublish = false\n\n[dependencies]\nlenso = {{ version = \"={}\", features = [\"host\"] }}\n{}\n",
seam.lenso_version,
dependencies.join("\n")
);
let source = format!(
"// {GENERATED_MARKER}; source-lock: {lock_digest}\n\npub const SOURCE_APPLICATION_LOCK_DIGEST: &str = \"{lock_digest}\";\n\n#[must_use]\npub fn linked_modules() -> Vec<lenso::host::HostLinkedModule> {{\n vec![\n{}\n ]\n}}\n",
bindings.join("\n")
);
let root = seam.generated_crate_path.trim_end_matches('/');
Ok(vec![
PlannedFile::generated(&format!("{root}/Cargo.toml"), &cargo),
PlannedFile::generated(&format!("{root}/src/lib.rs"), &source),
])
}
fn topological_linked_modules(
module_lock: &ApplicationModuleLock,
) -> Result<Vec<&crate::LockedModule>, LinkedWorkspaceError> {
let by_id = module_lock
.modules
.iter()
.map(|module| (module.module_id.as_str(), module))
.collect::<BTreeMap<_, _>>();
let mut remaining = by_id.keys().copied().collect::<BTreeSet<_>>();
let mut emitted = BTreeSet::new();
let mut ordered = Vec::new();
while !remaining.is_empty() {
let ready = remaining
.iter()
.copied()
.filter(|module_id| {
by_id[module_id]
.dependency_module_ids
.iter()
.all(|dependency| emitted.contains(dependency.as_str()))
})
.collect::<Vec<_>>();
if ready.is_empty() {
return invalid("Application Module Lock contains a dependency cycle");
}
for module_id in ready {
remaining.remove(module_id);
emitted.insert(module_id);
ordered.push(by_id[module_id]);
}
}
Ok(ordered)
}
fn validate_seam(root: &Path, seam: &LinkedCompositionSeam) -> Result<(), LinkedWorkspaceError> {
if seam.protocol != LINKED_COMPOSITION_SEAM_PROTOCOL
|| semver::Version::parse(&seam.lenso_version).is_err()
|| !valid_package_name(&seam.dependency_name)
{
return invalid("unsupported or malformed Linked composition seam");
}
for path in [
&seam.host_manifest_path,
&seam.host_source_path,
&seam.generated_crate_path,
] {
normalized_relative(Path::new(path))?;
}
let manifest = read_required_utf8(root, &seam.host_manifest_path)?;
let dependency_path = format!("path = \"{}\"", seam.generated_crate_path);
if !manifest.contains(&seam.dependency_name) || !manifest.contains(&dependency_path) {
return invalid(
"Host manifest does not contain the fixed generated composition dependency",
);
}
let source = read_required_utf8(root, &seam.host_source_path)?;
let rust_name = seam.dependency_name.replace('-', "_");
let call = format!(".linked_modules({rust_name}::linked_modules())");
if !source.contains(&call) {
return invalid("Host source does not call the fixed generated composition seam");
}
Ok(())
}
fn verify_read_set(
root: &Path,
plan: &crate::ModuleChangePlan,
effects: &[&ModulePlanEffect],
allow_applied_targets: bool,
) -> Result<(), LinkedWorkspaceError> {
for expected in &plan.read_set {
let snapshot = snapshot_path(root, &expected.path)?;
let observed = &snapshot.precondition;
if observed == expected {
if let Some(effect) = effects.iter().find(|effect| {
matches!(effect, ModulePlanEffect::WorkspaceFile { path, .. } if path == &expected.path)
}) {
verify_generated_ownership(plan, effect, snapshot.bytes.as_deref())?;
}
continue;
}
let applied = allow_applied_targets
&& effects.iter().any(|effect| {
workspace_effect_parts(effect).is_ok_and(|(path, change, after, _, _)| {
path == expected.path && matches_after(observed, change, after)
})
});
if !applied {
return Err(LinkedWorkspaceError::Stale(expected.path.clone()));
}
}
Ok(())
}
fn verify_generated_ownership(
plan: &crate::ModuleChangePlan,
effect: &ModulePlanEffect,
current_bytes: Option<&[u8]>,
) -> Result<(), LinkedWorkspaceError> {
let ModulePlanEffect::WorkspaceFile {
path,
ownership,
change,
..
} = effect
else {
return Ok(());
};
if *ownership != ModuleFileOwnership::Generated || *change == ModuleFileChange::Create {
return Ok(());
}
let bytes = current_bytes.ok_or_else(|| LinkedWorkspaceError::Stale(path.clone()))?;
let recognized = if path == "lenso.modules.lock.json" {
serde_json::from_slice::<ApplicationModuleLock>(bytes)
.ok()
.and_then(|module_lock| application_module_lock_digest(&module_lock).ok())
.as_deref()
== plan.current_lock_digest.as_deref()
} else {
let content = String::from_utf8_lossy(bytes);
content.contains(GENERATED_MARKER)
&& plan
.current_lock_digest
.as_deref()
.is_some_and(|digest| content.contains(digest))
};
if recognized {
Ok(())
} else {
Err(LinkedWorkspaceError::InvalidContract(format!(
"generated file `{path}` has no matching ownership marker and source lock"
)))
}
}
fn collect_backups(
root: &Path,
effects: &[&ModulePlanEffect],
) -> Result<Vec<ModuleWorkspaceBackup>, LinkedWorkspaceError> {
let mut paths = BTreeSet::new();
for effect in effects {
let (path, _, _, _, _) = workspace_effect_parts(effect)?;
let mut current = Path::new(path);
loop {
paths.insert(normalized_relative(current)?);
let Some(parent) = current.parent() else {
break;
};
if parent.as_os_str().is_empty() {
break;
}
current = parent;
}
}
paths
.into_iter()
.map(|path| {
let snapshot = snapshot_path(root, &path)?;
Ok(ModuleWorkspaceBackup {
path,
existence: snapshot.precondition.existence,
file_type: snapshot.precondition.file_type,
content_base64: snapshot.bytes.as_ref().map(|bytes| BASE64.encode(bytes)),
content_digest: snapshot.precondition.content_digest,
mode: snapshot.precondition.mode,
})
})
.collect()
}
fn restore_backups(
root: &Path,
backups: &[ModuleWorkspaceBackup],
) -> Result<(), LinkedWorkspaceError> {
for backup in backups.iter().rev() {
let path = guarded_path(root, &backup.path)?;
match (backup.existence, backup.file_type) {
(PathExistence::Present, ManagedFileType::Regular) => {
let bytes = BASE64
.decode(backup.content_base64.as_deref().ok_or_else(|| {
LinkedWorkspaceError::InvalidContract("file backup has no bytes".into())
})?)
.map_err(|error| LinkedWorkspaceError::InvalidContract(error.to_string()))?;
atomic_replace(&path, &bytes, backup.mode, "restore")?;
}
(PathExistence::Present, ManagedFileType::Directory) => {
fs::create_dir_all(&path)?;
set_mode(&path, backup.mode)?;
}
(PathExistence::Absent, _) => {
if path.is_file() {
fs::remove_file(&path)?;
} else if path.is_dir() {
fs::remove_dir(&path)?;
}
}
_ => return invalid("workspace backup has an incoherent shape"),
}
}
Ok(())
}
fn apply_file(
root: &Path,
path: &str,
change: ModuleFileChange,
after_content: Option<&str>,
after_mode: Option<u32>,
operation: &ModuleOperation,
effect_id: &str,
) -> Result<(), LinkedWorkspaceError> {
let target = guarded_path(root, path)?;
match change {
ModuleFileChange::Create | ModuleFileChange::Modify => {
let content = after_content.ok_or_else(|| {
LinkedWorkspaceError::InvalidContract("write effect has no reviewed bytes".into())
})?;
if let Some(parent) = target.parent() {
create_guarded_directories(root, parent)?;
}
atomic_replace(
&target,
content.as_bytes(),
after_mode,
&raw_digest(format!("{}:{effect_id}", operation.operation_id).as_bytes()),
)?;
}
ModuleFileChange::Delete => fs::remove_file(target)?,
}
Ok(())
}
fn workspace_effects(plan: &crate::ModuleChangePlan) -> Vec<&ModulePlanEffect> {
plan.effects
.iter()
.filter(|effect| matches!(effect, ModulePlanEffect::WorkspaceFile { .. }))
.collect()
}
#[allow(clippy::type_complexity)]
fn workspace_effect_parts(
effect: &ModulePlanEffect,
) -> Result<
(
&str,
ModuleFileChange,
Option<&str>,
Option<&str>,
Option<u32>,
),
LinkedWorkspaceError,
> {
if let ModulePlanEffect::WorkspaceFile {
path,
change,
after_digest,
after_content,
after_mode,
..
} = effect
{
Ok((
path,
*change,
after_digest.as_deref(),
after_content.as_deref(),
*after_mode,
))
} else {
invalid("non-workspace effect reached Linked workspace transaction")
}
}
fn matches_after(
observed: &ModulePathPrecondition,
change: ModuleFileChange,
after_digest: Option<&str>,
) -> bool {
match change {
ModuleFileChange::Create | ModuleFileChange::Modify => {
observed.existence == PathExistence::Present
&& observed.file_type == ManagedFileType::Regular
&& observed.content_digest.as_deref() == after_digest
}
ModuleFileChange::Delete => observed.existence == PathExistence::Absent,
}
}
fn snapshot_path(root: &Path, relative: &str) -> Result<PathSnapshot, LinkedWorkspaceError> {
let path = guarded_path(root, relative)?;
let metadata = match fs::symlink_metadata(&path) {
Ok(metadata) => metadata,
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
return Ok(PathSnapshot {
precondition: ModulePathPrecondition {
path: relative.to_owned(),
existence: PathExistence::Absent,
content_digest: None,
file_type: ManagedFileType::Absent,
mode: None,
},
bytes: None,
});
}
Err(error) => return Err(error.into()),
};
if metadata.file_type().is_symlink() {
return Err(LinkedWorkspaceError::UnsafePath(relative.to_owned()));
}
let (file_type, bytes, digest) = if metadata.is_file() {
let bytes = fs::read(&path)?;
let digest = raw_digest(&bytes);
(ManagedFileType::Regular, Some(bytes), Some(digest))
} else if metadata.is_dir() {
(ManagedFileType::Directory, None, None)
} else {
return Err(LinkedWorkspaceError::UnsafePath(relative.to_owned()));
};
Ok(PathSnapshot {
precondition: ModulePathPrecondition {
path: relative.to_owned(),
existence: PathExistence::Present,
content_digest: digest,
file_type,
mode: mode(&metadata),
},
bytes,
})
}
fn guarded_path(root: &Path, relative: &str) -> Result<PathBuf, LinkedWorkspaceError> {
let relative = normalized_relative(Path::new(relative))?;
let mut current = root.to_path_buf();
for component in Path::new(&relative).components() {
let Component::Normal(segment) = component else {
return Err(LinkedWorkspaceError::UnsafePath(relative));
};
current.push(segment);
match fs::symlink_metadata(¤t) {
Ok(metadata) if metadata.file_type().is_symlink() => {
return Err(LinkedWorkspaceError::UnsafePath(relative));
}
Ok(_) => {}
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
Err(error) => return Err(error.into()),
}
}
Ok(current)
}
fn normalized_relative(path: &Path) -> Result<String, LinkedWorkspaceError> {
if path.as_os_str().is_empty() || path.is_absolute() {
return Err(LinkedWorkspaceError::UnsafePath(path.display().to_string()));
}
let mut parts = Vec::new();
for component in path.components() {
match component {
Component::Normal(part) => parts.push(part.to_string_lossy().into_owned()),
_ => return Err(LinkedWorkspaceError::UnsafePath(path.display().to_string())),
}
}
let value = parts.join("/");
if value == ".env" || value.ends_with("/.env") {
return Err(LinkedWorkspaceError::UnsafePath(value));
}
Ok(value)
}
fn create_guarded_directories(root: &Path, parent: &Path) -> Result<(), LinkedWorkspaceError> {
let relative = parent
.strip_prefix(root)
.map_err(|_| LinkedWorkspaceError::UnsafePath(parent.display().to_string()))?;
let mut current = root.to_path_buf();
for component in relative.components() {
let Component::Normal(segment) = component else {
return Err(LinkedWorkspaceError::UnsafePath(
parent.display().to_string(),
));
};
current.push(segment);
if current.exists() {
if fs::symlink_metadata(¤t)?.file_type().is_symlink() || !current.is_dir() {
return Err(LinkedWorkspaceError::UnsafePath(
current.display().to_string(),
));
}
} else {
fs::create_dir(¤t)?;
}
}
Ok(())
}
fn atomic_replace(
path: &Path,
bytes: &[u8],
mode: Option<u32>,
nonce: &str,
) -> Result<(), LinkedWorkspaceError> {
let parent = path
.parent()
.ok_or_else(|| LinkedWorkspaceError::UnsafePath(path.display().to_string()))?;
let suffix = raw_digest(nonce.as_bytes()).replace("sha256:", "");
let temporary = parent.join(format!(".lenso-next-{}", &suffix[..16]));
let created = match OpenOptions::new()
.write(true)
.create_new(true)
.open(&temporary)
{
Ok(mut file) => {
let write_result = file
.write_all(bytes)
.and_then(|()| file.sync_all())
.and_then(|()| set_mode(&temporary, mode));
if let Err(error) = write_result {
let _ = fs::remove_file(&temporary);
return Err(error.into());
}
true
}
Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => {
let metadata = fs::symlink_metadata(&temporary)?;
if metadata.file_type().is_symlink()
|| !metadata.is_file()
|| fs::read(&temporary)? != bytes
{
return Err(LinkedWorkspaceError::Stale(temporary.display().to_string()));
}
false
}
Err(error) => return Err(error.into()),
};
if let Err(error) = fs::rename(&temporary, path) {
if created {
let _ = fs::remove_file(&temporary);
}
return Err(error.into());
}
File::open(parent)?.sync_all()?;
Ok(())
}
fn read_required_utf8(root: &Path, path: &str) -> Result<String, LinkedWorkspaceError> {
let snapshot = snapshot_path(root, path)?;
let bytes = snapshot.bytes.ok_or_else(|| {
LinkedWorkspaceError::InvalidContract(format!("required seam file `{path}` is absent"))
})?;
String::from_utf8(bytes).map_err(|_| {
LinkedWorkspaceError::InvalidContract(format!("required seam file `{path}` is not UTF-8"))
})
}
fn dependency_alias(module_id: &str) -> String {
format!(
"lenso_module_{}",
module_id
.chars()
.map(|character| if character.is_ascii_alphanumeric() {
character
} else {
'_'
})
.collect::<String>()
)
}
fn valid_package_name(value: &str) -> bool {
!value.is_empty()
&& value
.bytes()
.all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_'))
}
fn valid_binding_path(value: &str) -> bool {
!value.is_empty()
&& value.split("::").all(|segment| {
!segment.is_empty()
&& segment
.bytes()
.all(|byte| byte.is_ascii_alphanumeric() || byte == b'_')
&& segment
.bytes()
.next()
.is_some_and(|byte| byte.is_ascii_alphabetic() || byte == b'_')
})
}
fn escape_toml(value: &str) -> String {
value.replace('\\', "\\\\").replace('"', "\\\"")
}
fn exact_replacement_patch(path: &str, before: Option<&str>, after: &str) -> String {
let mut rendered_patch = format!("--- a/{path}\n+++ b/{path}\n");
if let Some(before) = before {
for line in before.lines() {
rendered_patch.push('-');
rendered_patch.push_str(line);
rendered_patch.push('\n');
}
}
for line in after.lines() {
rendered_patch.push('+');
rendered_patch.push_str(line);
rendered_patch.push('\n');
}
rendered_patch
}
fn raw_digest(bytes: &[u8]) -> String {
const HEX: &[u8; 16] = b"0123456789abcdef";
let digest = Sha256::digest(bytes);
let mut hex = String::with_capacity(digest.len() * 2);
for byte in digest {
hex.push(char::from(HEX[usize::from(byte >> 4)]));
hex.push(char::from(HEX[usize::from(byte & 0x0f)]));
}
format!("sha256:{hex}")
}
#[cfg(unix)]
#[allow(clippy::unnecessary_wraps)]
fn mode(metadata: &fs::Metadata) -> Option<u32> {
use std::os::unix::fs::PermissionsExt as _;
Some(metadata.permissions().mode() & 0o777)
}
#[cfg(not(unix))]
fn mode(_: &fs::Metadata) -> Option<u32> {
None
}
#[cfg(unix)]
fn set_mode(path: &Path, mode: Option<u32>) -> Result<(), std::io::Error> {
use std::os::unix::fs::PermissionsExt as _;
if let Some(mode) = mode {
fs::set_permissions(path, fs::Permissions::from_mode(mode))?;
}
Ok(())
}
#[cfg(not(unix))]
fn set_mode(_: &Path, _: Option<u32>) -> Result<(), std::io::Error> {
Ok(())
}
fn invalid<T>(message: impl Into<String>) -> Result<T, LinkedWorkspaceError> {
Err(LinkedWorkspaceError::InvalidContract(message.into()))
}