use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256};
use crate::{CapabilityOperationKind, ResolvedAppPlan};
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct TransitionValidationError {
detail: String,
}
impl std::fmt::Display for TransitionValidationError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(&self.detail)
}
}
impl std::error::Error for TransitionValidationError {}
fn invalid(detail: impl Into<String>) -> TransitionValidationError {
TransitionValidationError {
detail: detail.into(),
}
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(deny_unknown_fields)]
pub struct PlanTransition {
schema_version: u32,
predecessor_digest: String,
successor_digest: String,
replaced_instances: Vec<String>,
}
impl PlanTransition {
pub fn between(
previous: &ResolvedAppPlan,
successor: &ResolvedAppPlan,
) -> Result<Self, TransitionValidationError> {
previous
.validate()
.map_err(|_| invalid("invalid predecessor snapshot"))?;
successor
.validate()
.map_err(|_| invalid("invalid successor snapshot"))?;
if previous.execution_lanes().len() != 1
|| previous.execution_lanes() != successor.execution_lanes()
|| previous.capability_bindings() != successor.capability_bindings()
|| previous.terminal_policy() != successor.terminal_policy()
|| previous.plugin_instances().len() != successor.plugin_instances().len()
{
return Err(invalid(
"transition changes topology, bindings, lanes or terminal policy",
));
}
let mut replaced_instances = Vec::new();
for (before, after) in previous
.plugin_instances()
.iter()
.zip(successor.plugin_instances())
{
if before.provided_capabilities().iter().any(|endpoint| {
endpoint.operations().iter().any(|operation| {
endpoint.operation_kind(operation) != Some(CapabilityOperationKind::Request)
})
}) {
return Err(invalid(
"minimal transitions require a Request-only snapshot",
));
}
let mut normalized = after.clone().with_configuration(before.configuration());
if before.package_revision() != after.package_revision() {
if before.execution_class().as_str() != "lenso.bun-process@1" {
return Err(invalid(
"only Bun artifact revisions can change without a Host relink",
));
}
normalized = normalized.with_package_revision(before.package_revision());
}
if normalized != *before {
return Err(invalid(
"transition changes an Instance interface or execution selection",
));
}
if before != after {
replaced_instances.push(before.instance_key().to_owned());
}
}
if replaced_instances.is_empty() {
return Err(invalid("transition has no replacement Instances"));
}
if previous.capability_bindings().iter().any(|binding| {
replaced_instances
.iter()
.any(|key| key == binding.consumer_instance())
&& replaced_instances
.iter()
.any(|key| key == binding.provider_instance())
}) {
return Err(invalid(
"minimal staging cannot replace both ends of a dependency",
));
}
Ok(Self {
schema_version: 1,
predecessor_digest: Self::snapshot_digest(previous)?,
successor_digest: Self::snapshot_digest(successor)?,
replaced_instances,
})
}
pub fn snapshot_digest(plan: &ResolvedAppPlan) -> Result<String, TransitionValidationError> {
let bytes =
serde_json::to_vec(plan).map_err(|_| invalid("cannot encode snapshot identity"))?;
Ok(format!("sha256:{:x}", Sha256::digest(bytes)))
}
pub fn predecessor_digest(&self) -> &str {
&self.predecessor_digest
}
pub fn successor_digest(&self) -> &str {
&self.successor_digest
}
pub fn replaced_instances(&self) -> &[String] {
&self.replaced_instances
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::{CapabilityEndpointPlan, ExecutionClassId, PluginInstancePlan};
fn snapshot(instance: PluginInstancePlan) -> ResolvedAppPlan {
ResolvedAppPlan::new(vec![instance], vec![])
}
#[test]
fn config_changes_bind_exact_digest_and_decode_without_granting_authority() {
let instance = PluginInstancePlan::new("one", "example.plugin");
let previous = snapshot(instance.clone());
let successor = snapshot(instance.with_configuration("new"));
let transition = PlanTransition::between(&previous, &successor).unwrap();
assert_eq!(transition.replaced_instances(), ["one"]);
assert_ne!(
transition.predecessor_digest(),
transition.successor_digest()
);
assert_eq!(
serde_json::from_slice::<PlanTransition>(&serde_json::to_vec(&transition).unwrap())
.unwrap(),
transition
);
assert_eq!(
PlanTransition::snapshot_digest(&previous).unwrap(),
transition.predecessor_digest()
);
}
#[test]
fn revision_replacement_is_bun_only_and_other_execution_fields_stay_exact() {
let instance =
PluginInstancePlan::new("one", "example.plugin").with_package_revision("old");
assert!(
PlanTransition::between(
&snapshot(instance.clone()),
&snapshot(instance.clone().with_package_revision("new"))
)
.is_err()
);
let bun = instance.with_execution_class(ExecutionClassId::bun_child_process());
assert!(
PlanTransition::between(
&snapshot(bun.clone()),
&snapshot(bun.clone().with_package_revision("new"))
)
.is_ok()
);
assert!(
PlanTransition::between(
&snapshot(bun.clone()),
&snapshot(bun.with_entrypoint("other"))
)
.is_err()
);
}
#[test]
fn no_op_membership_and_non_request_snapshots_need_generation_fallback() {
let instance = PluginInstancePlan::new("one", "example.plugin");
let previous = snapshot(instance.clone());
assert!(PlanTransition::between(&previous, &previous).is_err());
assert!(PlanTransition::between(&previous, &ResolvedAppPlan::empty()).is_err());
for endpoint in [
CapabilityEndpointPlan::new("example.events@1", "1.0.0", ["publish"])
.with_event_operation("publish"),
CapabilityEndpointPlan::new("example.stream@1", "1.0.0", ["open"])
.with_stream_operation("open"),
] {
let instance = instance.clone().with_capability(endpoint);
assert!(
PlanTransition::between(
&snapshot(instance.clone()),
&snapshot(instance.with_configuration("new"))
)
.is_err()
);
}
}
}