pub mod identity;
use std::{
collections::{BTreeMap, BTreeSet},
env, fs,
path::{Path, PathBuf},
};
use anyhow::{Context, bail};
use lenso_app_plan::ExecutionClassId;
use lenso_app_plan::authoring::{
HostCatalog, PluginDescriptor, PluginInstanceId, PluginRootInstance, PluginRootSnapshot,
ResolvedApp, resolve_plugin_root,
};
use lenso_plugin_bundle::{
ImplementationPolicy, VerifiedBundle, read_bundle_manifest, resolve_implementation,
verify_bundle_directory,
};
use crate::identity::{
classify_existing_plugin_id, validate_plugin_id_v1, validate_release_version,
};
mod configuration_authority;
mod selection_authority;
pub use configuration_authority::{
LocalPluginRootAuthority, PluginConfigurationApplication, PluginConfigurationAuthority,
PluginConfigurationAuthoritySource, PluginConfigurationDiagnostic, PluginConfigurationProposal,
PluginConfigurationProposalStatus, PluginConfigurationPublication,
PluginConfigurationSourceConflict, PluginConfigurationSourceDigest, PluginRootRevision,
PluginRootRevisionConflict, PluginRootRevisionParseError, propose_instance_configuration,
publish_instance_configuration,
};
pub use selection_authority::{
PluginSelectionAuthority, PluginSelectionPublication, set_instance_enabled_fenced,
};
const PLUGIN_ROOT: &str = "plugins";
const HOST_CATALOG: &str = ".lenso/host-catalog.json";
const BUNDLE_NAME: &str = "plugin.lenso-plugin";
const AUTHORING_LOCK: &str = ".lenso/plugin-root-authoring.lock";
const MAX_CONFIGURATION_BYTES: u64 = 256 * 1024;
const MAX_RESOURCE_FILES: usize = 4_096;
const MAX_RESOURCE_FILE_BYTES: u64 = 1024 * 1024;
const MAX_RESOURCE_TOTAL_BYTES: u64 = 16 * 1024 * 1024;
const MAX_RESOURCE_DEPTH: usize = 32;
pub fn load_resolved_app(root: &Path) -> anyhow::Result<ResolvedApp> {
let host = load_host_catalog(root)?;
let snapshot = snapshot_plugin_root(root)?;
resolve_plugin_root(&host, &snapshot).map_err(anyhow::Error::msg)
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct PluginInstanceAuthoringState {
id: PluginInstanceId,
origin: PluginInstanceOrigin,
selection: PluginInstanceSelection,
root_configuration_toml: Option<String>,
source_digest: PluginConfigurationSourceDigest,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum PluginInstanceOrigin {
HostDefault { disableable: bool },
PluginRoot,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum PluginInstanceSelection {
Enabled,
DisabledByRoot,
}
impl PluginInstanceAuthoringState {
pub const fn id(&self) -> &PluginInstanceId {
&self.id
}
pub const fn is_enabled(&self) -> bool {
matches!(self.selection, PluginInstanceSelection::Enabled)
}
pub const fn is_host_default(&self) -> bool {
matches!(self.origin, PluginInstanceOrigin::HostDefault { .. })
}
pub const fn is_disableable(&self) -> bool {
match self.origin {
PluginInstanceOrigin::HostDefault { disableable } => disableable,
PluginInstanceOrigin::PluginRoot => true,
}
}
pub fn root_configuration_toml(&self) -> Option<&str> {
self.root_configuration_toml.as_deref()
}
pub const fn source_digest(&self) -> &PluginConfigurationSourceDigest {
&self.source_digest
}
pub const fn is_disabled_by_root(&self) -> bool {
matches!(self.selection, PluginInstanceSelection::DisabledByRoot)
}
pub const fn has_root_difference(&self) -> bool {
self.root_configuration_toml.is_some() || self.is_disabled_by_root()
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct PluginAuthoringState {
plugin_id: String,
release_version: String,
root_supplied: bool,
instances: Vec<PluginInstanceAuthoringState>,
}
impl PluginAuthoringState {
pub fn plugin_id(&self) -> &str {
&self.plugin_id
}
pub fn release_version(&self) -> &str {
&self.release_version
}
pub const fn is_root_supplied(&self) -> bool {
self.root_supplied
}
pub fn instances(&self) -> &[PluginInstanceAuthoringState] {
&self.instances
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct PluginRootAuthoringState {
revision: PluginRootRevision,
resolved: ResolvedApp,
plugins: Vec<PluginAuthoringState>,
}
impl PluginRootAuthoringState {
pub const fn revision(&self) -> &PluginRootRevision {
&self.revision
}
pub const fn resolved(&self) -> &ResolvedApp {
&self.resolved
}
pub fn plugins(&self) -> &[PluginAuthoringState] {
&self.plugins
}
}
pub fn inspect_plugin_root(root: &Path) -> anyhow::Result<PluginRootAuthoringState> {
let host = load_host_catalog(root)?;
let snapshot = snapshot_plugin_root(root)?;
let revision = configuration_authority::revision_for_snapshot(&snapshot)?;
let resolved = resolve_plugin_root(&host, &snapshot).map_err(anyhow::Error::msg)?;
let enabled = resolved
.instances()
.iter()
.map(|instance| instance.id().clone())
.collect::<BTreeSet<_>>();
let disabled = snapshot.disabled().iter().cloned().collect::<BTreeSet<_>>();
let root_instances = snapshot
.instances()
.iter()
.map(|instance| instance.id().clone())
.collect::<BTreeSet<_>>();
let host_defaults = host
.defaults()
.iter()
.map(|instance| (instance.id().clone(), instance.is_disableable()))
.collect::<BTreeMap<_, _>>();
let ids = root_instances
.iter()
.chain(disabled.iter())
.chain(host_defaults.keys())
.cloned()
.collect::<BTreeSet<_>>();
let root_releases = snapshot
.releases()
.iter()
.map(|release| release.plugin_id().to_owned())
.collect::<BTreeSet<_>>();
let mut releases = host
.plugins()
.iter()
.map(|release| {
(
release.descriptor().plugin_id().to_owned(),
release.descriptor().release_version().to_owned(),
)
})
.chain(snapshot.releases().iter().map(|release| {
(
release.plugin_id().to_owned(),
release.release_version().to_owned(),
)
}))
.collect::<BTreeMap<_, _>>();
for id in &ids {
releases
.entry(id.plugin_id().to_owned())
.or_insert_with(String::new);
}
let mut plugins = Vec::with_capacity(releases.len());
for (plugin_id, release_version) in releases {
let plugin_ids = ids
.iter()
.filter(|id| id.plugin_id() == plugin_id)
.cloned()
.collect::<Vec<_>>();
let mut instances = Vec::with_capacity(plugin_ids.len());
for id in plugin_ids {
let configuration_path = root
.join(PLUGIN_ROOT)
.join(id.plugin_id())
.join(format!("{}.toml", id.instance_key()));
let root_configuration_toml = if root_instances.contains(&id) {
Some(fs::read_to_string(&configuration_path).with_context(|| {
format!(
"read Plugin configuration source {}",
configuration_path.display()
)
})?)
} else {
None
};
let source_digest = instance_source_digest(&id, root_configuration_toml.as_deref());
let host_disableable = host_defaults.get(&id).copied();
instances.push(PluginInstanceAuthoringState {
origin: host_disableable.map_or(PluginInstanceOrigin::PluginRoot, |disableable| {
PluginInstanceOrigin::HostDefault { disableable }
}),
selection: if enabled.contains(&id) {
PluginInstanceSelection::Enabled
} else {
PluginInstanceSelection::DisabledByRoot
},
root_configuration_toml,
source_digest,
id,
});
}
plugins.push(PluginAuthoringState {
root_supplied: root_releases.contains(&plugin_id),
plugin_id,
release_version,
instances,
});
}
Ok(authoring_state(revision, resolved, plugins))
}
fn instance_source_digest(
id: &PluginInstanceId,
source: Option<&str>,
) -> PluginConfigurationSourceDigest {
configuration_authority::source_digest_for_bytes(
id.plugin_id(),
id.instance_key(),
source.map(str::as_bytes),
)
}
fn authoring_state(
revision: PluginRootRevision,
resolved: ResolvedApp,
plugins: Vec<PluginAuthoringState>,
) -> PluginRootAuthoringState {
PluginRootAuthoringState {
revision,
resolved,
plugins,
}
}
fn load_host_catalog(root: &Path) -> anyhow::Result<HostCatalog> {
let path = root.join(HOST_CATALOG);
let metadata = fs::symlink_metadata(&path).with_context(|| {
format!(
"Host Catalog is unavailable at {}; build or install the current Host first",
path.display()
)
})?;
if !metadata.file_type().is_file() {
bail!("Host Catalog must be a regular file: {}", path.display());
}
let bytes = fs::read(&path).with_context(|| format!("read Host Catalog {}", path.display()))?;
serde_json::from_slice(&bytes)
.with_context(|| format!("Host Catalog is invalid: {}", path.display()))
}
fn snapshot_plugin_root(root: &Path) -> anyhow::Result<PluginRootSnapshot> {
let plugin_root = root.join(PLUGIN_ROOT);
match fs::symlink_metadata(&plugin_root) {
Ok(metadata) if metadata.file_type().is_dir() => {}
Ok(_) => bail!(
"Plugin Root must be a regular directory: {}",
plugin_root.display()
),
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
return Ok(PluginRootSnapshot::default());
}
Err(error) => {
return Err(error).with_context(|| format!("inspect {}", plugin_root.display()));
}
}
let mut releases = Vec::new();
let mut instances = Vec::new();
let mut disabled = Vec::new();
let mut plugin_names = BTreeMap::<String, String>::new();
let mut directories = read_entries(&plugin_root)?;
directories.sort_by_key(fs::DirEntry::file_name);
for entry in directories {
let name = utf8_name(&entry.path(), &entry.file_name())?;
if is_ignored_os_metadata(&name) {
continue;
}
let file_type = entry.file_type()?;
if !file_type.is_dir() {
bail!("unknown Plugin Root entry: {}", entry.path().display());
}
let plugin_id = name;
validate_existing_plugin_id(&plugin_id)?;
reject_case_collision(&mut plugin_names, &plugin_id, "Plugin ID")?;
scan_plugin_directory(
&entry.path(),
&plugin_id,
&mut releases,
&mut instances,
&mut disabled,
)?;
}
Ok(PluginRootSnapshot::new(releases, instances, disabled))
}
fn scan_plugin_directory(
directory: &Path,
plugin_id: &str,
releases: &mut Vec<PluginDescriptor>,
instances: &mut Vec<PluginRootInstance>,
disabled: &mut Vec<PluginInstanceId>,
) -> anyhow::Result<()> {
let mut normalized = BTreeMap::<String, String>::new();
let mut configured_instances = BTreeSet::new();
let mut resource_directories = BTreeMap::<String, PathBuf>::new();
let mut entries = read_entries(directory)?;
entries.sort_by_key(fs::DirEntry::file_name);
for entry in entries {
let name = utf8_name(&entry.path(), &entry.file_name())?;
if is_ignored_os_metadata(&name) {
continue;
}
reject_case_collision(&mut normalized, &name, "Plugin filename")?;
let file_type = entry.file_type()?;
if name == BUNDLE_NAME {
if !file_type.is_dir() {
bail!(
"Plugin Bundle must be a regular directory: {}",
entry.path().display()
);
}
releases.push(read_bundle_descriptor(&entry.path(), plugin_id)?);
continue;
}
if file_type.is_dir() {
validate_instance_filename(&name)?;
resource_directories.insert(name, entry.path());
continue;
}
if !file_type.is_file() {
bail!(
"Plugin entries cannot be symlinks or special files: {}",
entry.path().display()
);
}
if let Some(instance) = name.strip_suffix(".toml") {
validate_instance_filename(instance)?;
configured_instances.insert(instance.to_owned());
instances.push(
PluginRootInstance::new(plugin_id, instance)
.with_configuration(read_configuration(&entry.path())?),
);
} else if let Some(instance) = name.strip_suffix(".disabled") {
validate_instance_filename(instance)?;
if fs::metadata(entry.path())?.len() != 0 {
bail!("disabled marker must be empty: {}", entry.path().display());
}
disabled.push(PluginInstanceId::new(plugin_id, instance));
} else {
bail!("unknown Plugin file: {}", entry.path().display());
}
}
for (instance, resource_directory) in resource_directories {
if !configured_instances.contains(&instance) {
bail!(
"orphan Plugin resource directory without `{instance}.toml`: {}",
resource_directory.display()
);
}
validate_resource_directory(&resource_directory)?;
}
Ok(())
}
fn validate_resource_directory(path: &Path) -> anyhow::Result<()> {
let mut file_count = 0_usize;
let mut total_size = 0_u64;
let mut pending = vec![(path.to_path_buf(), 0_usize)];
while let Some((directory, depth)) = pending.pop() {
if depth > MAX_RESOURCE_DEPTH {
bail!(
"Plugin resource directory exceeds {MAX_RESOURCE_DEPTH} levels: {}",
directory.display()
);
}
let mut entries = read_entries(&directory)?;
entries.sort_by_key(fs::DirEntry::file_name);
for entry in entries {
let entry_path = entry.path();
let name = utf8_name(&entry_path, &entry.file_name())?;
if is_ignored_os_metadata(&name) {
continue;
}
let file_type = entry.file_type()?;
if file_type.is_dir() {
pending.push((entry_path, depth + 1));
continue;
}
if !file_type.is_file() {
bail!(
"Plugin resources cannot contain symlinks or special files: {}",
entry_path.display()
);
}
if file_count == MAX_RESOURCE_FILES {
bail!(
"Plugin resources exceed {MAX_RESOURCE_FILES} files: {}",
path.display()
);
}
let metadata = fs::symlink_metadata(&entry_path)?;
if metadata.file_type().is_symlink() || !metadata.file_type().is_file() {
bail!(
"Plugin resources must be regular files: {}",
entry_path.display()
);
}
if metadata.len() > MAX_RESOURCE_FILE_BYTES {
bail!("Plugin resource exceeds 1 MiB: {}", entry_path.display());
}
let bytes = fs::read(&entry_path)?;
let byte_count = u64::try_from(bytes.len()).with_context(|| {
format!("Plugin resource is too large: {}", entry_path.display())
})?;
if byte_count > MAX_RESOURCE_FILE_BYTES {
bail!("Plugin resource exceeds 1 MiB: {}", entry_path.display());
}
total_size = total_size
.checked_add(byte_count)
.with_context(|| format!("Plugin resource size overflow: {}", path.display()))?;
if total_size > MAX_RESOURCE_TOTAL_BYTES {
bail!("Plugin resources exceed 16 MiB: {}", path.display());
}
file_count += 1;
}
}
Ok(())
}
fn is_ignored_os_metadata(name: &str) -> bool {
name == ".DS_Store"
}
fn read_bundle_descriptor(path: &Path, plugin_id: &str) -> anyhow::Result<PluginDescriptor> {
validate_existing_plugin_id(plugin_id)?;
let verified = verify_bundle_directory(path)
.with_context(|| format!("verify Plugin Bundle {}", path.display()))?;
read_verified_bundle_descriptor(path, plugin_id, &verified)
}
fn read_verified_bundle_descriptor(
path: &Path,
plugin_id: &str,
verified: &VerifiedBundle,
) -> anyhow::Result<PluginDescriptor> {
if verified.plugin_id != plugin_id {
bail!(
"Plugin Bundle ID `{}` does not match directory `{plugin_id}`",
verified.plugin_id
);
}
let manifest = read_bundle_manifest(path)
.with_context(|| format!("read Plugin Manifest {}", path.display()))?;
let descriptor = resolve_implementation(
&manifest,
&ImplementationPolicy {
host_target: format!("{}-unknown-{}", env::consts::ARCH, env::consts::OS),
execution_classes: vec![
ExecutionClassId::new("lenso.quickjs@1"),
ExecutionClassId::new("lenso.process@1"),
ExecutionClassId::new("lenso.wasm-component@1"),
ExecutionClassId::new("lenso.bun-process@1"),
],
},
)?
.descriptor;
if descriptor.plugin_id() != plugin_id
|| descriptor.release_version() != verified.release_version
{
bail!("Plugin Descriptor identity does not match the verified Bundle");
}
Ok(descriptor)
}
fn read_configuration(path: &Path) -> anyhow::Result<serde_json::Value> {
let metadata = fs::metadata(path)?;
if metadata.len() > MAX_CONFIGURATION_BYTES {
bail!("Plugin configuration exceeds 256 KiB: {}", path.display());
}
let text = fs::read_to_string(path)
.with_context(|| format!("read Plugin configuration {}", path.display()))?;
let table: toml::Table = toml::from_str(&text)
.with_context(|| format!("parse Plugin configuration {}", path.display()))?;
serde_json::to_value(table).context("convert Plugin configuration to portable values")
}
fn read_entries(path: &Path) -> anyhow::Result<Vec<fs::DirEntry>> {
fs::read_dir(path)
.with_context(|| format!("read directory {}", path.display()))?
.collect::<Result<Vec<_>, _>>()
.with_context(|| format!("read directory entries {}", path.display()))
}
fn utf8_name(path: &Path, name: &std::ffi::OsStr) -> anyhow::Result<String> {
name.to_str()
.map(str::to_owned)
.with_context(|| format!("Plugin path is not UTF-8: {}", path.display()))
}
fn validate_instance_filename(instance: &str) -> anyhow::Result<()> {
validate_path_identity(instance, "Instance key")?;
if instance.starts_with('.') || instance == "plugin" {
bail!("reserved Plugin Instance key `{instance}`");
}
Ok(())
}
fn validate_existing_plugin_id(plugin_id: &str) -> anyhow::Result<()> {
validate_path_identity(plugin_id, "Plugin ID")?;
classify_existing_plugin_id(plugin_id).map(|_| ())
}
fn validate_path_identity(value: &str, label: &str) -> anyhow::Result<()> {
if value.trim() != value
|| value.is_empty()
|| value == "."
|| value == ".."
|| value.contains(['/', '\0', '\\'])
{
bail!("invalid {label} `{value}`");
}
Ok(())
}
fn reject_case_collision(
normalized: &mut BTreeMap<String, String>,
value: &str,
label: &str,
) -> anyhow::Result<()> {
let key = value.to_lowercase();
if let Some(previous) = normalized.insert(key, value.to_owned())
&& previous != value
{
bail!("case-colliding {label}s `{previous}` and `{value}`");
}
Ok(())
}
pub fn add_bundle(root: &Path, bundle: &Path) -> anyhow::Result<(String, String, ResolvedApp)> {
prepare_bundle_mutation(root, bundle, BundleMutation::Add)?.commit()
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum BundleMutation {
Add,
Replace,
Restore,
}
#[derive(Debug)]
pub struct PreparedBundleMutation {
authority: fs::File,
destination: PathBuf,
mutation: BundleMutation,
resolved: ResolvedApp,
staging: tempfile::TempDir,
verified: VerifiedBundle,
}
impl PreparedBundleMutation {
pub const fn verified(&self) -> &VerifiedBundle {
&self.verified
}
pub const fn resolved(&self) -> &ResolvedApp {
&self.resolved
}
pub fn destination(&self) -> &Path {
&self.destination
}
pub fn commit(self) -> anyhow::Result<(String, String, ResolvedApp)> {
let Self {
authority,
destination,
mutation,
resolved,
staging,
verified,
} = self;
let commit = commit_staged_bundle(&destination, mutation, staging);
drop(authority);
commit?;
Ok((verified.plugin_id, verified.release_version, resolved))
}
}
fn commit_staged_bundle(
destination: &Path,
mutation: BundleMutation,
staging: tempfile::TempDir,
) -> anyhow::Result<()> {
commit_staged_bundle_with(
destination,
mutation,
staging,
atomic_publish_bundle,
tempfile::TempDir::close,
)
}
fn commit_staged_bundle_with<Publish, Retire>(
destination: &Path,
mutation: BundleMutation,
staging: tempfile::TempDir,
publish: Publish,
retire: Retire,
) -> anyhow::Result<()>
where
Publish: FnOnce(&Path, &Path, BundleMutation) -> std::io::Result<()>,
Retire: FnOnce(tempfile::TempDir) -> std::io::Result<()>,
{
let parent = destination
.parent()
.context("Bundle destination has no parent")?;
if mutation == BundleMutation::Add && destination.exists() {
bail!("Plugin Bundle already exists: {}", destination.display());
}
let created_parent = mutation == BundleMutation::Add && !parent.exists();
if mutation == BundleMutation::Add {
fs::create_dir_all(parent)?;
}
let publication =
publish(staging.path(), destination, mutation).with_context(|| match mutation {
BundleMutation::Add => format!("commit Plugin Bundle {}", destination.display()),
BundleMutation::Replace | BundleMutation::Restore => {
format!("atomically replace Plugin Bundle {}", destination.display())
}
});
if let Err(error) = publication {
if created_parent
&& let Err(cleanup_error) = fs::remove_dir(parent)
&& cleanup_error.kind() != std::io::ErrorKind::NotFound
&& cleanup_error.kind() != std::io::ErrorKind::DirectoryNotEmpty
{
return Err(error.context(format!(
"also failed to remove empty Plugin directory {}: {cleanup_error}",
parent.display()
)));
}
return Err(error);
}
if mutation != BundleMutation::Add
&& let Err(error) = retire(staging)
{
eprintln!("warning: Plugin Bundle committed, but retired Bundle cleanup failed: {error}");
}
Ok(())
}
#[cfg(any(target_os = "linux", target_vendor = "apple"))]
fn atomic_publish_bundle(
staging: &Path,
destination: &Path,
mutation: BundleMutation,
) -> std::io::Result<()> {
use rustix::fs::{CWD, RenameFlags, renameat_with};
let flags = match mutation {
BundleMutation::Add => RenameFlags::NOREPLACE,
BundleMutation::Replace | BundleMutation::Restore => RenameFlags::EXCHANGE,
};
renameat_with(CWD, staging, CWD, destination, flags).map_err(std::io::Error::from)
}
#[cfg(windows)]
fn atomic_publish_bundle(
staging: &Path,
destination: &Path,
mutation: BundleMutation,
) -> std::io::Result<()> {
match mutation {
BundleMutation::Add => winsafe::MoveFile(
staging.to_str().ok_or_else(|| {
std::io::Error::new(
std::io::ErrorKind::InvalidInput,
"Plugin Bundle staging path is not Unicode",
)
})?,
destination.to_str().ok_or_else(|| {
std::io::Error::new(
std::io::ErrorKind::InvalidInput,
"Plugin Bundle destination path is not Unicode",
)
})?,
)
.map_err(|error| std::io::Error::from_raw_os_error(error.raw() as i32)),
BundleMutation::Replace | BundleMutation::Restore => Err(std::io::Error::new(
std::io::ErrorKind::Unsupported,
"atomic Plugin Bundle replacement is unavailable on this platform",
)),
}
}
#[cfg(not(any(target_os = "linux", target_vendor = "apple", windows)))]
fn atomic_publish_bundle(
_staging: &Path,
_destination: &Path,
_mutation: BundleMutation,
) -> std::io::Result<()> {
Err(std::io::Error::new(
std::io::ErrorKind::Unsupported,
"atomic Plugin Bundle publication is unavailable on this platform",
))
}
pub fn prepare_bundle_mutation(
root: &Path,
bundle: &Path,
mutation: BundleMutation,
) -> anyhow::Result<PreparedBundleMutation> {
let staging = tempfile::Builder::new()
.prefix(".plugin-bundle-")
.tempdir_in(root)?;
copy_directory(bundle, staging.path())?;
let (verified, descriptor) = verify_bundle_mutation(staging.path(), mutation)?;
let authority = lock_plugin_root(root)?;
let resolved = resolve_bundle_mutation(root, mutation, &verified, descriptor)?;
let destination = root
.join(PLUGIN_ROOT)
.join(&verified.plugin_id)
.join(BUNDLE_NAME);
Ok(PreparedBundleMutation {
authority,
destination,
mutation,
resolved,
staging,
verified,
})
}
pub fn validate_bundle_mutation(
root: &Path,
bundle: &Path,
mutation: BundleMutation,
) -> anyhow::Result<(lenso_plugin_bundle::VerifiedBundle, ResolvedApp)> {
let (verified, descriptor) = verify_bundle_mutation(bundle, mutation)?;
let _lock = lock_plugin_root(root)?;
let resolved = resolve_bundle_mutation(root, mutation, &verified, descriptor)?;
Ok((verified, resolved))
}
fn verify_bundle_mutation(
bundle: &Path,
mutation: BundleMutation,
) -> anyhow::Result<(VerifiedBundle, PluginDescriptor)> {
let verified = verify_bundle_directory(bundle)
.with_context(|| format!("verify Plugin Bundle {}", bundle.display()))?;
match mutation {
BundleMutation::Add | BundleMutation::Replace => {
validate_plugin_id_v1(&verified.plugin_id)?;
}
BundleMutation::Restore => {
classify_existing_plugin_id(&verified.plugin_id)?;
}
}
validate_release_version(&verified.release_version)?;
let descriptor = read_verified_bundle_descriptor(bundle, &verified.plugin_id, &verified)?;
Ok((verified, descriptor))
}
fn resolve_bundle_mutation(
root: &Path,
mutation: BundleMutation,
verified: &VerifiedBundle,
descriptor: PluginDescriptor,
) -> anyhow::Result<ResolvedApp> {
let host = load_host_catalog(root)?;
let current = snapshot_plugin_root(root)?;
let has_current = current
.releases()
.iter()
.any(|release| release.plugin_id() == verified.plugin_id);
match (mutation, has_current) {
(BundleMutation::Add, true) => {
bail!("Plugin `{}` already has a root Bundle", verified.plugin_id)
}
(BundleMutation::Replace | BundleMutation::Restore, false) => {
bail!(
"Plugin `{}` has no root Bundle to update",
verified.plugin_id
)
}
_ => {}
}
let candidate = PluginRootSnapshot::new(
current
.releases()
.iter()
.filter(|release| release.plugin_id() != verified.plugin_id)
.cloned()
.chain([descriptor]),
current.instances().iter().cloned(),
current.disabled().iter().cloned(),
);
let resolved = resolve_plugin_root(&host, &candidate).map_err(anyhow::Error::msg)?;
Ok(resolved)
}
pub fn configure_instance(
root: &Path,
plugin_id: &str,
instance: &str,
bytes: &[u8],
) -> anyhow::Result<ResolvedApp> {
let base_revision = inspect_plugin_root(root)?.revision().clone();
let proposal =
propose_instance_configuration(root, &base_revision, plugin_id, instance, bytes)?;
let publication = publish_instance_configuration(root, &proposal)?;
Ok(publication.into_resolved())
}
pub fn set_instance_disabled(
root: &Path,
plugin_id: &str,
instance: &str,
disabled_state: bool,
) -> anyhow::Result<ResolvedApp> {
set_instance_disabled_inner(root, plugin_id, instance, disabled_state, None)
.map(|(_, _, resolved)| resolved)
}
fn set_instance_disabled_inner(
root: &Path,
plugin_id: &str,
instance: &str,
disabled_state: bool,
expected_revision: Option<&PluginRootRevision>,
) -> anyhow::Result<(PluginRootRevision, PluginRootRevision, ResolvedApp)> {
validate_existing_plugin_id(plugin_id)?;
validate_instance_filename(instance)?;
let _lock = lock_plugin_root(root)?;
let host = load_host_catalog(root)?;
let current = snapshot_plugin_root(root)?;
let base_revision = configuration_authority::revision_for_snapshot(¤t)?;
if let Some(expected_revision) = expected_revision {
configuration_authority::ensure_revision(expected_revision, &base_revision)?;
}
let id = PluginInstanceId::new(plugin_id, instance);
let mut disabled = current.disabled().iter().cloned().collect::<BTreeSet<_>>();
if disabled_state {
disabled.insert(id.clone());
} else if !disabled.remove(&id) {
bail!("Plugin Instance `{id}` is not disabled");
}
let candidate = PluginRootSnapshot::new(
current.releases().iter().cloned(),
current.instances().iter().cloned(),
disabled,
);
let candidate_revision = configuration_authority::revision_for_snapshot(&candidate)?;
let resolved = resolve_plugin_root(&host, &candidate).map_err(anyhow::Error::msg)?;
let marker = root
.join(PLUGIN_ROOT)
.join(plugin_id)
.join(format!("{instance}.disabled"));
if disabled_state {
atomic_write(&marker, &[])?;
} else {
fs::remove_file(&marker)
.with_context(|| format!("remove disabled marker {}", marker.display()))?;
}
Ok((base_revision, candidate_revision, resolved))
}
pub fn remove_instance_difference(
root: &Path,
plugin_id: &str,
instance: &str,
) -> anyhow::Result<ResolvedApp> {
validate_existing_plugin_id(plugin_id)?;
validate_instance_filename(instance)?;
let _lock = lock_plugin_root(root)?;
let host = load_host_catalog(root)?;
let current = snapshot_plugin_root(root)?;
let id = PluginInstanceId::new(plugin_id, instance);
let candidate = PluginRootSnapshot::new(
current.releases().iter().cloned(),
current
.instances()
.iter()
.filter(|item| item.id() != &id)
.cloned(),
current
.disabled()
.iter()
.filter(|item| *item != &id)
.cloned(),
);
let resolved = resolve_plugin_root(&host, &candidate).map_err(anyhow::Error::msg)?;
let plugin_directory = root.join(PLUGIN_ROOT).join(plugin_id);
remove_if_exists(&plugin_directory.join(format!("{instance}.toml")))?;
remove_if_exists(&plugin_directory.join(format!("{instance}.disabled")))?;
Ok(resolved)
}
pub fn remove_plugin(root: &Path, plugin_id: &str) -> anyhow::Result<(ResolvedApp, PathBuf)> {
validate_existing_plugin_id(plugin_id)?;
let _lock = lock_plugin_root(root)?;
let host = load_host_catalog(root)?;
let current = snapshot_plugin_root(root)?;
let candidate = PluginRootSnapshot::new(
current
.releases()
.iter()
.filter(|release| release.plugin_id() != plugin_id)
.cloned(),
current
.instances()
.iter()
.filter(|instance| instance.id().plugin_id() != plugin_id)
.cloned(),
current
.disabled()
.iter()
.filter(|instance| instance.plugin_id() != plugin_id)
.cloned(),
);
let resolved = resolve_plugin_root(&host, &candidate).map_err(anyhow::Error::msg)?;
let plugin_directory = root.join(PLUGIN_ROOT).join(plugin_id);
if !plugin_directory.exists() {
bail!("Plugin `{plugin_id}` has no Plugin Root directory");
}
let trash = root
.join(".lenso/trash")
.join(format!("{plugin_id}-{}", uuid::Uuid::now_v7()));
fs::create_dir_all(trash.parent().expect("trash has a parent"))?;
fs::rename(&plugin_directory, &trash)?;
Ok((resolved, trash))
}
fn atomic_write(path: &Path, bytes: &[u8]) -> anyhow::Result<()> {
let parent = path.parent().context("Plugin file has no parent")?;
fs::create_dir_all(parent)?;
let temporary = tempfile::NamedTempFile::new_in(parent)?;
fs::write(temporary.path(), bytes)?;
temporary
.persist(path)
.map_err(|error| error.error)
.with_context(|| format!("commit Plugin file {}", path.display()))?;
Ok(())
}
fn lock_plugin_root(root: &Path) -> anyhow::Result<fs::File> {
let path = root.join(AUTHORING_LOCK);
let parent = path.parent().context("Plugin Root lock has no parent")?;
fs::create_dir_all(parent)?;
let file = fs::OpenOptions::new()
.create(true)
.read(true)
.write(true)
.truncate(false)
.open(&path)
.with_context(|| format!("open Plugin Root authoring lock {}", path.display()))?;
file.lock()
.with_context(|| format!("lock Plugin Root authoring authority {}", path.display()))?;
Ok(file)
}
fn copy_directory(source: &Path, destination: &Path) -> anyhow::Result<()> {
for entry in read_entries(source)? {
let file_type = entry.file_type()?;
if file_type.is_dir() {
let child = destination.join(entry.file_name());
fs::create_dir_all(&child)?;
copy_directory(&entry.path(), &child)?;
continue;
}
if !file_type.is_file() {
bail!(
"Plugin Bundle contains a non-file entry: {}",
entry.path().display()
);
}
fs::copy(entry.path(), destination.join(entry.file_name()))?;
}
Ok(())
}
fn remove_if_exists(path: &Path) -> anyhow::Result<()> {
match fs::remove_file(path) {
Ok(()) => Ok(()),
Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
Err(error) => Err(error).with_context(|| format!("remove {}", path.display())),
}
}
#[cfg(test)]
mod tests {
use super::*;
use lenso_app_plan::authoring::{HostDefaultPlugin, HostPluginRelease, HostSlot};
fn fixture_root() -> tempfile::TempDir {
let root = tempfile::tempdir().unwrap();
fs::create_dir_all(root.path().join(".lenso")).unwrap();
let host = HostCatalog::new(
[HostSlot::one("agent")],
[HostPluginRelease::new(PluginDescriptor::new(
"example.agent",
"1.0.0",
"agent",
))],
[HostDefaultPlugin::new("example.agent", "default")],
);
fs::write(
root.path().join(HOST_CATALOG),
serde_json::to_vec(&host).unwrap(),
)
.unwrap();
root
}
#[test]
fn missing_plugin_root_resolves_the_host_default_app() {
let root = fixture_root();
let resolved = load_resolved_app(root.path()).unwrap();
assert_eq!(resolved.instances().len(), 1);
assert_eq!(
resolved.instances()[0].id().to_string(),
"example.agent/default"
);
}
#[test]
fn inspection_separates_host_defaults_from_root_differences() {
let root = fixture_root();
let plugin = root.path().join("plugins/example.agent");
fs::create_dir_all(&plugin).unwrap();
fs::write(plugin.join("default.toml"), "").unwrap();
let state = inspect_plugin_root(root.path()).unwrap();
let plugin = state
.plugins()
.iter()
.find(|plugin| plugin.plugin_id() == "example.agent")
.unwrap();
let instance = &plugin.instances()[0];
assert_eq!(plugin.release_version(), "1.0.0");
assert!(!plugin.is_root_supplied());
assert!(instance.is_enabled());
assert!(instance.is_host_default());
assert!(!instance.is_disableable());
assert_eq!(instance.root_configuration_toml(), Some(""));
assert!(instance.source_digest().as_str().starts_with("sha256:"));
assert!(instance.has_root_difference());
}
#[test]
fn inspection_reports_disabled_host_default_without_losing_the_instance() {
let root = tempfile::tempdir().unwrap();
fs::create_dir_all(root.path().join(".lenso")).unwrap();
let host = HostCatalog::new(
[HostSlot::optional("optional")],
[HostPluginRelease::new(PluginDescriptor::new(
"example.optional",
"1.0.0",
"optional",
))],
[HostDefaultPlugin::new("example.optional", "default").disableable()],
);
fs::write(
root.path().join(HOST_CATALOG),
serde_json::to_vec(&host).unwrap(),
)
.unwrap();
let plugin = root.path().join("plugins/example.optional");
fs::create_dir_all(&plugin).unwrap();
fs::write(plugin.join("default.disabled"), "").unwrap();
let state = inspect_plugin_root(root.path()).unwrap();
let instance = &state.plugins()[0].instances()[0];
assert!(!instance.is_enabled());
assert!(instance.is_host_default());
assert!(instance.is_disableable());
assert!(instance.is_disabled_by_root());
}
#[test]
fn local_selection_authority_disables_and_enables_one_instance() {
let root = tempfile::tempdir().unwrap();
fs::create_dir_all(root.path().join(".lenso")).unwrap();
let host = HostCatalog::new(
[HostSlot::optional("optional")],
[HostPluginRelease::new(PluginDescriptor::new(
"example.optional",
"1.0.0",
"optional",
))],
[HostDefaultPlugin::new("example.optional", "default").disableable()],
);
fs::write(
root.path().join(HOST_CATALOG),
serde_json::to_vec(&host).unwrap(),
)
.unwrap();
let authority = LocalPluginRootAuthority::new(root.path());
let base = inspect_plugin_root(root.path()).unwrap().revision().clone();
let disabled = authority
.set_enabled(&base, "example.optional", "default", false)
.unwrap();
assert_eq!(disabled.base_revision(), &base);
assert!(!disabled.enabled());
assert_eq!(disabled.plugin_id(), "example.optional");
assert_eq!(disabled.instance(), "default");
assert!(
root.path()
.join("plugins/example.optional/default.disabled")
.is_file()
);
let enabled = authority
.set_enabled(disabled.revision(), "example.optional", "default", true)
.unwrap();
assert!(enabled.enabled());
assert!(
!root
.path()
.join("plugins/example.optional/default.disabled")
.exists()
);
}
#[test]
fn local_selection_authority_rejects_a_stale_revision_without_mutating() {
let root = tempfile::tempdir().unwrap();
fs::create_dir_all(root.path().join(".lenso")).unwrap();
let host = HostCatalog::new(
[HostSlot::optional("optional")],
[HostPluginRelease::new(PluginDescriptor::new(
"example.optional",
"1.0.0",
"optional",
))],
[HostDefaultPlugin::new("example.optional", "default").disableable()],
);
fs::write(
root.path().join(HOST_CATALOG),
serde_json::to_vec(&host).unwrap(),
)
.unwrap();
let authority = LocalPluginRootAuthority::new(root.path());
let stale = inspect_plugin_root(root.path()).unwrap().revision().clone();
authority
.set_enabled(&stale, "example.optional", "default", false)
.unwrap();
let error = authority
.set_enabled(&stale, "example.optional", "default", true)
.unwrap_err();
assert!(error.downcast_ref::<PluginRootRevisionConflict>().is_some());
assert!(
root.path()
.join("plugins/example.optional/default.disabled")
.is_file()
);
}
#[test]
fn macos_metadata_at_plugin_root_is_ignored() {
let root = fixture_root();
fs::create_dir(root.path().join("plugins")).unwrap();
fs::write(root.path().join("plugins/.DS_Store"), b"Finder metadata").unwrap();
let resolved = load_resolved_app(root.path()).unwrap();
assert_eq!(resolved.instances().len(), 1);
}
#[test]
fn macos_metadata_inside_plugin_directory_is_ignored() {
let root = fixture_root();
let plugin = root.path().join("plugins/example.agent");
fs::create_dir_all(&plugin).unwrap();
fs::write(plugin.join(".DS_Store"), b"Finder metadata").unwrap();
let resolved = load_resolved_app(root.path()).unwrap();
assert_eq!(resolved.instances().len(), 1);
}
#[test]
fn accepts_a_bounded_resource_directory_paired_with_an_instance() {
let root = fixture_root();
let plugin = root.path().join("plugins/example.agent");
fs::create_dir_all(plugin.join("default/prompts")).unwrap();
fs::write(plugin.join("default.toml"), "").unwrap();
fs::write(plugin.join("default/prompts/system.md"), "hello").unwrap();
fs::write(plugin.join("default/prompts/.DS_Store"), "metadata").unwrap();
let resolved = load_resolved_app(root.path()).unwrap();
assert!(
resolved
.instances()
.iter()
.any(|instance| instance.id().to_string() == "example.agent/default")
);
}
#[test]
fn rejects_an_orphan_resource_directory() {
let root = fixture_root();
let resources = root.path().join("plugins/example.agent/custom");
fs::create_dir_all(&resources).unwrap();
fs::write(resources.join("prompt.md"), "orphan").unwrap();
let error = load_resolved_app(root.path()).unwrap_err();
assert!(
error
.to_string()
.contains("orphan Plugin resource directory")
);
}
#[cfg(unix)]
#[test]
fn rejects_a_resource_symlink() {
use std::os::unix::fs::symlink;
let root = fixture_root();
let plugin = root.path().join("plugins/example.agent");
fs::create_dir_all(plugin.join("custom")).unwrap();
fs::write(plugin.join("custom.toml"), "").unwrap();
fs::write(root.path().join("secret"), "not admitted").unwrap();
symlink(root.path().join("secret"), plugin.join("custom/secret")).unwrap();
let error = load_resolved_app(root.path()).unwrap_err();
assert!(error.to_string().contains("cannot contain symlinks"));
}
#[test]
fn failed_configuration_candidate_does_not_write_the_plugin_root() {
let root = fixture_root();
let error = configure_instance(
root.path(),
"example.agent",
"default",
b"unexpected = true\n",
)
.unwrap_err();
assert!(error.to_string().contains("non-empty configuration"));
assert!(
!root
.path()
.join("plugins/example.agent/default.toml")
.exists()
);
}
#[test]
fn required_default_disable_fails_before_writing_a_marker() {
let root = fixture_root();
let error =
set_instance_disabled(root.path(), "example.agent", "default", true).unwrap_err();
assert!(error.to_string().contains("cannot be disabled"));
assert!(
!root
.path()
.join("plugins/example.agent/default.disabled")
.exists()
);
}
#[test]
fn case_colliding_plugin_identities_fail_closed() {
let mut normalized = BTreeMap::new();
reject_case_collision(&mut normalized, "Example.Agent", "Plugin ID").unwrap();
let error =
reject_case_collision(&mut normalized, "example.agent", "Plugin ID").unwrap_err();
assert!(error.to_string().contains("case-colliding Plugin IDs"));
}
#[test]
fn add_replace_and_restore_publish_failures_leave_visible_bytes_unchanged() {
for mutation in [
BundleMutation::Add,
BundleMutation::Replace,
BundleMutation::Restore,
] {
let root = tempfile::tempdir().unwrap();
let destination = root
.path()
.join("plugins/example.agent/plugin.lenso-plugin");
if mutation == BundleMutation::Add {
fs::create_dir(root.path().join("plugins")).unwrap();
} else {
fs::create_dir_all(&destination).unwrap();
fs::write(destination.join("marker"), "old").unwrap();
}
let staging = tempfile::tempdir_in(root.path()).unwrap();
fs::write(staging.path().join("marker"), "new").unwrap();
let error = commit_staged_bundle_with(
&destination,
mutation,
staging,
|_, _, _| {
Err(std::io::Error::new(
std::io::ErrorKind::PermissionDenied,
"injected publish failure",
))
},
|_| panic!("retirement cannot run before publication succeeds"),
)
.unwrap_err();
assert!(error.to_string().contains("Plugin Bundle"));
if mutation == BundleMutation::Add {
assert!(!destination.exists());
assert!(!destination.parent().unwrap().exists());
} else {
assert_eq!(
fs::read_to_string(destination.join("marker")).unwrap(),
"old"
);
}
}
}
#[test]
fn portable_bundle_add_publishes_with_one_atomic_rename() {
let root = tempfile::tempdir().unwrap();
let destination = root
.path()
.join("plugins/example.agent/plugin.lenso-plugin");
let staging = tempfile::tempdir_in(root.path()).unwrap();
fs::write(staging.path().join("marker"), "new").unwrap();
commit_staged_bundle(&destination, BundleMutation::Add, staging).unwrap();
assert_eq!(
fs::read_to_string(destination.join("marker")).unwrap(),
"new"
);
}
#[cfg(any(target_os = "linux", target_vendor = "apple", windows))]
#[test]
fn portable_bundle_add_never_replaces_a_concurrent_destination() {
let root = tempfile::tempdir().unwrap();
let destination = root.path().join("destination");
fs::create_dir(&destination).unwrap();
fs::write(destination.join("marker"), "old").unwrap();
let staging = tempfile::tempdir_in(root.path()).unwrap();
fs::write(staging.path().join("marker"), "new").unwrap();
atomic_publish_bundle(staging.path(), &destination, BundleMutation::Add).unwrap_err();
assert_eq!(
fs::read_to_string(destination.join("marker")).unwrap(),
"old"
);
assert_eq!(
fs::read_to_string(staging.path().join("marker")).unwrap(),
"new"
);
}
#[cfg(not(any(target_os = "linux", target_vendor = "apple")))]
#[test]
fn portable_bundle_replace_fails_closed_when_exchange_is_unavailable() {
let root = tempfile::tempdir().unwrap();
let destination = root
.path()
.join("plugins/example.agent/plugin.lenso-plugin");
fs::create_dir_all(&destination).unwrap();
fs::write(destination.join("marker"), "old").unwrap();
let staging = tempfile::tempdir_in(root.path()).unwrap();
fs::write(staging.path().join("marker"), "new").unwrap();
let error =
commit_staged_bundle(&destination, BundleMutation::Replace, staging).unwrap_err();
assert_eq!(
error
.root_cause()
.downcast_ref::<std::io::Error>()
.unwrap()
.kind(),
std::io::ErrorKind::Unsupported
);
assert_eq!(
fs::read_to_string(destination.join("marker")).unwrap(),
"old"
);
}
#[cfg(any(target_os = "linux", target_vendor = "apple"))]
#[test]
fn replace_and_restore_commit_atomically_even_when_retirement_cleanup_fails() {
for mutation in [BundleMutation::Replace, BundleMutation::Restore] {
let root = tempfile::tempdir().unwrap();
let destination = root
.path()
.join("plugins/example.agent/plugin.lenso-plugin");
fs::create_dir_all(&destination).unwrap();
fs::write(destination.join("marker"), "old").unwrap();
let staging = tempfile::tempdir_in(root.path()).unwrap();
fs::write(staging.path().join("marker"), "new").unwrap();
let mut retired = None;
commit_staged_bundle_with(
&destination,
mutation,
staging,
atomic_publish_bundle,
|staging| {
retired = Some(staging.keep());
Err(std::io::Error::other("injected cleanup failure"))
},
)
.unwrap();
assert_eq!(
fs::read_to_string(destination.join("marker")).unwrap(),
"new"
);
let retired = retired.unwrap();
assert_eq!(fs::read_to_string(retired.join("marker")).unwrap(), "old");
fs::remove_dir_all(retired).unwrap();
}
}
}