pub mod authoring;
use std::{
collections::{BTreeMap, BTreeSet},
fmt,
time::Duration,
};
use serde::{Deserialize, Serialize};
mod execution;
mod resolution;
pub use execution::{ExecutionClassId, ExecutionLaneId, ExecutionLanePlan};
use resolution::{
activation_order_for, resolve_parts, sort_bindings, sort_module_instances,
sorted_execution_lanes, validate_execution_lanes,
};
pub const PLAN_SCHEMA_VERSION: u32 = 1;
pub const DEFAULT_REQUEST_QUEUE_CAPACITY: usize = 16;
pub const DEFAULT_REQUEST_MAX_CONCURRENCY: usize = 1;
pub const DEFAULT_EVENT_QUEUE_CAPACITY: usize = 16;
fn default_execution_lanes() -> Vec<ExecutionLanePlan> {
vec![ExecutionLanePlan::new("main")]
}
#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum CapabilityCardinality {
One,
Optional,
Many,
}
#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum CapabilityOperationKind {
Request,
Stream,
Event,
}
#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct RequestAdmissionPlan {
queue_capacity: usize,
max_concurrency: usize,
}
impl RequestAdmissionPlan {
pub const fn new(queue_capacity: usize, max_concurrency: usize) -> Self {
Self {
queue_capacity,
max_concurrency,
}
}
pub const fn queue_capacity(self) -> usize {
self.queue_capacity
}
pub const fn max_concurrency(self) -> usize {
self.max_concurrency
}
fn validate(self, capability_id: &str, operation: &str) -> Result<(), PlanResolutionError> {
if self.max_concurrency == 0 {
return Err(PlanResolutionError::InvalidRequestAdmission {
capability_id: capability_id.to_owned(),
operation: operation.to_owned(),
queue_capacity: self.queue_capacity,
max_concurrency: self.max_concurrency,
});
}
Ok(())
}
}
impl Default for RequestAdmissionPlan {
fn default() -> Self {
Self::new(
DEFAULT_REQUEST_QUEUE_CAPACITY,
DEFAULT_REQUEST_MAX_CONCURRENCY,
)
}
}
#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct EventAdmissionPlan {
capacity: usize,
}
impl EventAdmissionPlan {
pub const fn new(capacity: usize) -> Self {
Self { capacity }
}
pub const fn capacity(self) -> usize {
self.capacity
}
}
impl Default for EventAdmissionPlan {
fn default() -> Self {
Self::new(DEFAULT_EVENT_QUEUE_CAPACITY)
}
}
#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum RestartMode {
Never,
OnFailure,
}
#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct RestartPolicy {
mode: RestartMode,
max_attempts: usize,
window: Duration,
backoff: Duration,
jitter: Duration,
stability: Duration,
}
impl RestartPolicy {
pub const fn never() -> Self {
Self {
mode: RestartMode::Never,
max_attempts: 0,
window: Duration::ZERO,
backoff: Duration::ZERO,
jitter: Duration::ZERO,
stability: Duration::ZERO,
}
}
pub const fn on_failure(
max_attempts: usize,
window: Duration,
backoff: Duration,
jitter: Duration,
stability: Duration,
) -> Self {
Self {
mode: RestartMode::OnFailure,
max_attempts,
window,
backoff,
jitter,
stability,
}
}
pub const fn mode(self) -> RestartMode {
self.mode
}
pub const fn max_attempts(self) -> usize {
self.max_attempts
}
pub const fn window(self) -> Duration {
self.window
}
pub const fn backoff(self) -> Duration {
self.backoff
}
pub const fn jitter(self) -> Duration {
self.jitter
}
pub const fn stability(self) -> Duration {
self.stability
}
fn validate(&self, instance_key: &str) -> Result<(), PlanResolutionError> {
if self.mode == RestartMode::OnFailure && (self.max_attempts == 0 || self.window.is_zero())
{
return Err(PlanResolutionError::InvalidRestartPolicy {
instance_key: instance_key.to_owned(),
max_attempts: self.max_attempts,
window: self.window,
});
}
Ok(())
}
}
impl Default for RestartPolicy {
fn default() -> Self {
Self::never()
}
}
#[derive(Clone, Copy, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum ModuleCriticality {
#[default]
NonCritical,
Critical,
}
impl ModuleCriticality {
pub const fn is_critical(self) -> bool {
matches!(self, Self::Critical)
}
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct CapabilityRequirementPlan {
capability_id: String,
descriptor_version: String,
cardinality: CapabilityCardinality,
}
impl CapabilityRequirementPlan {
pub fn new(
capability_id: impl Into<String>,
descriptor_version: impl Into<String>,
cardinality: CapabilityCardinality,
) -> Self {
Self {
capability_id: capability_id.into(),
descriptor_version: descriptor_version.into(),
cardinality,
}
}
pub fn one(capability_id: impl Into<String>, descriptor_version: impl Into<String>) -> Self {
Self::new(
capability_id,
descriptor_version,
CapabilityCardinality::One,
)
}
pub fn optional(
capability_id: impl Into<String>,
descriptor_version: impl Into<String>,
) -> Self {
Self::new(
capability_id,
descriptor_version,
CapabilityCardinality::Optional,
)
}
pub fn many(capability_id: impl Into<String>, descriptor_version: impl Into<String>) -> Self {
Self::new(
capability_id,
descriptor_version,
CapabilityCardinality::Many,
)
}
pub fn capability_id(&self) -> &str {
&self.capability_id
}
pub fn descriptor_version(&self) -> &str {
&self.descriptor_version
}
pub const fn cardinality(&self) -> CapabilityCardinality {
self.cardinality
}
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct CapabilityEndpointPlan {
capability_id: String,
descriptor_version: String,
operations: Vec<String>,
operation_kinds: BTreeMap<String, CapabilityOperationKind>,
default_admission: Option<RequestAdmissionPlan>,
operation_admissions: BTreeMap<String, RequestAdmissionPlan>,
event_admission: Option<EventAdmissionPlan>,
#[serde(default)]
cross_lane_transfer: bool,
}
impl CapabilityEndpointPlan {
pub fn new(
capability_id: impl Into<String>,
descriptor_version: impl Into<String>,
operations: impl IntoIterator<Item = impl Into<String>>,
) -> Self {
Self {
capability_id: capability_id.into(),
descriptor_version: descriptor_version.into(),
operations: operations.into_iter().map(Into::into).collect(),
operation_kinds: BTreeMap::new(),
default_admission: None,
operation_admissions: BTreeMap::new(),
event_admission: None,
cross_lane_transfer: false,
}
}
#[must_use]
pub fn with_admission(mut self, admission: RequestAdmissionPlan) -> Self {
self.default_admission = Some(admission);
self
}
#[must_use]
pub fn with_limits(self, queue_capacity: usize, max_concurrency: usize) -> Self {
self.with_admission(RequestAdmissionPlan::new(queue_capacity, max_concurrency))
}
#[must_use]
pub fn with_operation_kind(
mut self,
operation: impl Into<String>,
kind: CapabilityOperationKind,
) -> Self {
self.operation_kinds.insert(operation.into(), kind);
self
}
#[must_use]
pub fn with_stream_operation(self, operation: impl Into<String>) -> Self {
self.with_operation_kind(operation, CapabilityOperationKind::Stream)
}
#[must_use]
pub fn with_event_operation(self, operation: impl Into<String>) -> Self {
self.with_operation_kind(operation, CapabilityOperationKind::Event)
}
#[must_use]
pub fn with_event_admission(mut self, admission: EventAdmissionPlan) -> Self {
self.event_admission = Some(admission);
self
}
#[must_use]
pub fn with_event_capacity(self, capacity: usize) -> Self {
self.with_event_admission(EventAdmissionPlan::new(capacity))
}
#[must_use]
pub const fn with_cross_lane_transfer(mut self) -> Self {
self.cross_lane_transfer = true;
self
}
#[must_use]
pub fn with_operation_admission(
mut self,
operation: impl Into<String>,
admission: RequestAdmissionPlan,
) -> Self {
self.operation_admissions
.insert(operation.into(), admission);
self
}
#[must_use]
pub fn with_operation_limits(
self,
operation: impl Into<String>,
queue_capacity: usize,
max_concurrency: usize,
) -> Self {
self.with_operation_admission(
operation,
RequestAdmissionPlan::new(queue_capacity, max_concurrency),
)
}
pub fn capability_id(&self) -> &str {
&self.capability_id
}
pub fn descriptor_version(&self) -> &str {
&self.descriptor_version
}
pub fn operations(&self) -> &[String] {
&self.operations
}
pub fn operation_kind(&self, operation: &str) -> Option<CapabilityOperationKind> {
self.operations
.iter()
.any(|declared| declared == operation)
.then(|| {
self.operation_kinds
.get(operation)
.copied()
.unwrap_or(CapabilityOperationKind::Request)
})
}
pub fn stream_operations(&self) -> Vec<&str> {
self.operations
.iter()
.filter(|operation| {
self.operation_kind(operation) == Some(CapabilityOperationKind::Stream)
})
.map(String::as_str)
.collect()
}
pub fn request_operations(&self) -> Vec<&str> {
self.operations
.iter()
.filter(|operation| {
self.operation_kind(operation) == Some(CapabilityOperationKind::Request)
})
.map(String::as_str)
.collect()
}
pub fn event_operations(&self) -> Vec<&str> {
self.operations
.iter()
.filter(|operation| {
self.operation_kind(operation) == Some(CapabilityOperationKind::Event)
})
.map(String::as_str)
.collect()
}
pub const fn event_admission(&self) -> Option<EventAdmissionPlan> {
self.event_admission
}
pub const fn supports_cross_lane_transfer(&self) -> bool {
self.cross_lane_transfer
}
pub fn default_admission(&self) -> Option<RequestAdmissionPlan> {
self.default_admission
}
pub fn operation_admissions(&self) -> &BTreeMap<String, RequestAdmissionPlan> {
&self.operation_admissions
}
pub fn operation_admission(&self, operation: &str) -> Option<RequestAdmissionPlan> {
self.operation_admissions
.get(operation)
.copied()
.or(self.default_admission)
}
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct ModuleInstancePlan {
instance_key: String,
package_id: String,
entrypoint: String,
configuration: String,
provided_capabilities: Vec<CapabilityEndpointPlan>,
required_capabilities: Vec<CapabilityRequirementPlan>,
execution_class: ExecutionClassId,
package_revision: String,
restart_policy: RestartPolicy,
criticality: ModuleCriticality,
#[serde(default)]
execution_lane: ExecutionLaneId,
}
impl ModuleInstancePlan {
pub fn new(instance_key: impl Into<String>, package_id: impl Into<String>) -> Self {
Self {
instance_key: instance_key.into(),
package_id: package_id.into(),
entrypoint: "default".to_owned(),
configuration: "{}".to_owned(),
provided_capabilities: Vec::new(),
required_capabilities: Vec::new(),
execution_class: ExecutionClassId::native_rust(),
package_revision: String::new(),
restart_policy: RestartPolicy::default(),
criticality: ModuleCriticality::default(),
execution_lane: ExecutionLaneId::default(),
}
}
#[must_use]
pub fn with_entrypoint(mut self, entrypoint: impl Into<String>) -> Self {
self.entrypoint = entrypoint.into();
self
}
#[must_use]
pub fn with_configuration(mut self, configuration: impl Into<String>) -> Self {
self.configuration = configuration.into();
self
}
#[must_use]
pub fn with_capability(mut self, capability: CapabilityEndpointPlan) -> Self {
self.provided_capabilities.push(capability);
self
}
#[must_use]
pub fn with_requirement(mut self, requirement: CapabilityRequirementPlan) -> Self {
self.required_capabilities.push(requirement);
self
}
#[must_use]
pub fn with_required_capability(self, requirement: CapabilityRequirementPlan) -> Self {
self.with_requirement(requirement)
}
#[must_use]
pub fn with_execution_class(mut self, execution_class: ExecutionClassId) -> Self {
self.execution_class = execution_class;
self
}
#[must_use]
pub fn with_execution_lane(mut self, execution_lane: ExecutionLaneId) -> Self {
self.execution_lane = execution_lane;
self
}
#[must_use]
pub fn with_package_revision(mut self, revision: impl Into<String>) -> Self {
self.package_revision = revision.into();
self
}
#[must_use]
pub fn with_restart_policy(mut self, restart_policy: RestartPolicy) -> Self {
self.restart_policy = restart_policy;
self
}
#[must_use]
pub fn with_criticality(mut self, criticality: ModuleCriticality) -> Self {
self.criticality = criticality;
self
}
pub fn instance_key(&self) -> &str {
&self.instance_key
}
pub fn package_id(&self) -> &str {
&self.package_id
}
pub fn entrypoint(&self) -> &str {
&self.entrypoint
}
pub fn configuration(&self) -> &str {
&self.configuration
}
pub fn provided_capabilities(&self) -> &[CapabilityEndpointPlan] {
&self.provided_capabilities
}
pub fn required_capabilities(&self) -> &[CapabilityRequirementPlan] {
&self.required_capabilities
}
pub fn execution_class(&self) -> &ExecutionClassId {
&self.execution_class
}
pub const fn execution_lane(&self) -> &ExecutionLaneId {
&self.execution_lane
}
pub fn package_revision(&self) -> &str {
&self.package_revision
}
pub const fn restart_policy(&self) -> RestartPolicy {
self.restart_policy
}
pub const fn criticality(&self) -> ModuleCriticality {
self.criticality
}
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct CapabilityBinding {
consumer_instance: String,
capability_id: String,
descriptor_version: String,
provider_instance: String,
provider_order: usize,
admission: RequestAdmissionPlan,
admission_explicit: bool,
event_admission: EventAdmissionPlan,
event_admission_explicit: bool,
}
impl CapabilityBinding {
pub fn new(
consumer_instance: impl Into<String>,
capability_id: impl Into<String>,
descriptor_version: impl Into<String>,
provider_instance: impl Into<String>,
) -> Self {
Self {
consumer_instance: consumer_instance.into(),
capability_id: capability_id.into(),
descriptor_version: descriptor_version.into(),
provider_instance: provider_instance.into(),
provider_order: 0,
admission: RequestAdmissionPlan::default(),
admission_explicit: false,
event_admission: EventAdmissionPlan::default(),
event_admission_explicit: false,
}
}
#[must_use]
pub fn with_admission(mut self, admission: RequestAdmissionPlan) -> Self {
self.admission = admission;
self.admission_explicit = true;
self
}
#[must_use]
pub fn with_limits(self, queue_capacity: usize, max_concurrency: usize) -> Self {
self.with_admission(RequestAdmissionPlan::new(queue_capacity, max_concurrency))
}
#[must_use]
pub fn with_event_admission(mut self, admission: EventAdmissionPlan) -> Self {
self.event_admission = admission;
self.event_admission_explicit = true;
self
}
#[must_use]
pub fn with_event_capacity(self, capacity: usize) -> Self {
self.with_event_admission(EventAdmissionPlan::new(capacity))
}
fn with_provider_order(mut self, provider_order: usize) -> Self {
self.provider_order = provider_order;
self
}
pub fn consumer_instance(&self) -> &str {
&self.consumer_instance
}
pub fn capability_id(&self) -> &str {
&self.capability_id
}
pub fn descriptor_version(&self) -> &str {
&self.descriptor_version
}
pub fn provider_instance(&self) -> &str {
&self.provider_instance
}
pub const fn provider_order(&self) -> usize {
self.provider_order
}
pub const fn admission(&self) -> RequestAdmissionPlan {
self.admission
}
pub const fn has_explicit_admission(&self) -> bool {
self.admission_explicit
}
pub const fn event_admission(&self) -> EventAdmissionPlan {
self.event_admission
}
pub const fn has_explicit_event_admission(&self) -> bool {
self.event_admission_explicit
}
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct AppComposition {
module_instances: Vec<ModuleInstancePlan>,
capability_bindings: Vec<CapabilityBinding>,
#[serde(default = "default_execution_lanes")]
execution_lanes: Vec<ExecutionLanePlan>,
}
impl AppComposition {
pub fn new(
module_instances: Vec<ModuleInstancePlan>,
capability_bindings: Vec<CapabilityBinding>,
) -> Self {
Self {
module_instances,
capability_bindings,
execution_lanes: default_execution_lanes(),
}
}
#[must_use]
pub fn with_execution_lanes(mut self, execution_lanes: Vec<ExecutionLanePlan>) -> Self {
self.execution_lanes = execution_lanes;
self
}
pub fn resolve(&self) -> Result<ResolvedAppPlan, PlanResolutionError> {
validate_execution_lanes(&self.execution_lanes, &self.module_instances)?;
resolve_parts(&self.module_instances, &self.capability_bindings).map(
|(module_instances, capability_bindings)| ResolvedAppPlan {
schema_version: PLAN_SCHEMA_VERSION,
module_instances,
capability_bindings,
execution_lanes: sorted_execution_lanes(&self.execution_lanes),
},
)
}
pub fn module_instances(&self) -> &[ModuleInstancePlan] {
&self.module_instances
}
pub fn capability_bindings(&self) -> &[CapabilityBinding] {
&self.capability_bindings
}
pub fn execution_lanes(&self) -> &[ExecutionLanePlan] {
&self.execution_lanes
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub enum PlanResolutionError {
UnsupportedSchemaVersion { expected: u32, actual: u32 },
DuplicateModuleInstance { instance_key: String },
MissingExecutionLane,
InvalidExecutionLane { execution_lane: String },
DuplicateExecutionLane { execution_lane: String },
UndeclaredExecutionLane {
instance_key: String,
execution_lane: String,
},
InvalidModuleEntrypoint { instance_key: String },
DuplicateProvidedCapability {
provider_instance: String,
capability_id: String,
},
DuplicateOperation {
provider_instance: String,
capability_id: String,
operation: String,
},
DuplicateRequiredCapability {
consumer_instance: String,
capability_id: String,
},
InvalidConsumerReference {
consumer_instance: String,
capability_id: String,
},
InvalidProviderReference {
consumer_instance: String,
capability_id: String,
provider_instance: String,
},
UndeclaredCapabilityRequirement {
consumer_instance: String,
capability_id: String,
},
IncompatibleCapabilityVersion {
consumer_instance: String,
capability_id: String,
required: String,
provided: String,
provider_instance: String,
},
CrossLaneTransferUnsupported {
consumer_instance: String,
provider_instance: String,
capability_id: String,
},
CrossLaneInteractionUnsupported {
capability_id: String,
operation: String,
interaction: CapabilityOperationKind,
},
MissingOneBinding {
consumer_instance: String,
capability_id: String,
},
AmbiguousOneBinding {
consumer_instance: String,
capability_id: String,
providers: usize,
},
AmbiguousOptionalBinding {
consumer_instance: String,
capability_id: String,
providers: usize,
},
DuplicateBinding {
consumer_instance: String,
capability_id: String,
provider_instance: String,
},
InvalidRequestAdmission {
capability_id: String,
operation: String,
queue_capacity: usize,
max_concurrency: usize,
},
UnknownAdmissionOperation {
capability_id: String,
operation: String,
},
UnknownOperationInteraction {
capability_id: String,
operation: String,
},
InvalidRestartPolicy {
instance_key: String,
max_attempts: usize,
window: Duration,
},
ActivationCycle { instances: Vec<String> },
}
impl fmt::Display for PlanResolutionError {
#[allow(clippy::too_many_lines)]
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::UnsupportedSchemaVersion { expected, actual } => write!(
formatter,
"unsupported Plan schema version {actual}; expected {expected}"
),
Self::DuplicateModuleInstance { instance_key } => {
write!(formatter, "duplicate Module Instance `{instance_key}`")
}
Self::MissingExecutionLane => {
formatter.write_str("Resolved App Plan declares no Execution Lanes")
}
Self::InvalidExecutionLane { execution_lane } => {
write!(formatter, "invalid Execution Lane `{execution_lane}`")
}
Self::DuplicateExecutionLane { execution_lane } => {
write!(formatter, "duplicate Execution Lane `{execution_lane}`")
}
Self::UndeclaredExecutionLane {
instance_key,
execution_lane,
} => write!(
formatter,
"Module Instance `{instance_key}` is placed on undeclared Execution Lane `{execution_lane}`"
),
Self::InvalidModuleEntrypoint { instance_key } => write!(
formatter,
"Module Instance `{instance_key}` has an empty entrypoint"
),
Self::DuplicateProvidedCapability {
provider_instance,
capability_id,
} => write!(
formatter,
"Module Instance `{provider_instance}` provides Capability `{capability_id}` more than once"
),
Self::DuplicateOperation {
provider_instance,
capability_id,
operation,
} => write!(
formatter,
"Module Instance `{provider_instance}` Capability `{capability_id}` declares Operation `{operation}` more than once"
),
Self::DuplicateRequiredCapability {
consumer_instance,
capability_id,
} => write!(
formatter,
"Module Instance `{consumer_instance}` requires Capability `{capability_id}` more than once"
),
Self::InvalidConsumerReference {
consumer_instance,
capability_id,
} => write!(
formatter,
"Capability `{capability_id}` names missing consumer `{consumer_instance}`"
),
Self::InvalidProviderReference {
consumer_instance,
capability_id,
provider_instance,
} => write!(
formatter,
"consumer `{consumer_instance}` names invalid provider `{provider_instance}` for Capability `{capability_id}`"
),
Self::UndeclaredCapabilityRequirement {
consumer_instance,
capability_id,
} => write!(
formatter,
"consumer `{consumer_instance}` has no declared requirement for Capability `{capability_id}`"
),
Self::IncompatibleCapabilityVersion {
consumer_instance,
capability_id,
required,
provided,
provider_instance,
} => write!(
formatter,
"consumer `{consumer_instance}` requires Capability `{capability_id}` version `{required}`, but provider `{provider_instance}` provides `{provided}`"
),
Self::CrossLaneTransferUnsupported {
consumer_instance,
provider_instance,
capability_id,
} => write!(
formatter,
"consumer `{consumer_instance}` binds Capability `{capability_id}` across Execution Lanes to provider `{provider_instance}`, but its contract types do not support cross-lane transfer"
),
Self::CrossLaneInteractionUnsupported {
capability_id,
operation,
interaction,
} => write!(
formatter,
"Capability `{capability_id}` Operation `{operation}` uses {interaction:?}, which this Plan version cannot transfer across Execution Lanes"
),
Self::MissingOneBinding {
consumer_instance,
capability_id,
} => write!(
formatter,
"consumer `{consumer_instance}` is missing one binding for Capability `{capability_id}`"
),
Self::AmbiguousOneBinding {
consumer_instance,
capability_id,
providers,
} => write!(
formatter,
"consumer `{consumer_instance}` has {providers} bindings for one Capability `{capability_id}`"
),
Self::AmbiguousOptionalBinding {
consumer_instance,
capability_id,
providers,
} => write!(
formatter,
"consumer `{consumer_instance}` has {providers} bindings for optional Capability `{capability_id}`"
),
Self::DuplicateBinding {
consumer_instance,
capability_id,
provider_instance,
} => write!(
formatter,
"consumer `{consumer_instance}` binds Capability `{capability_id}` to provider `{provider_instance}` more than once"
),
Self::InvalidRequestAdmission {
capability_id,
operation,
queue_capacity,
max_concurrency,
} => write!(
formatter,
"Capability `{capability_id}` Operation `{operation}` has invalid request admission (queue capacity {queue_capacity}, concurrency {max_concurrency})"
),
Self::UnknownAdmissionOperation {
capability_id,
operation,
} => write!(
formatter,
"Capability `{capability_id}` configures request admission for unknown Operation `{operation}`"
),
Self::UnknownOperationInteraction {
capability_id,
operation,
} => write!(
formatter,
"Capability `{capability_id}` configures interaction metadata for unknown Operation `{operation}`"
),
Self::InvalidRestartPolicy {
instance_key,
max_attempts,
window,
} => write!(
formatter,
"Module Instance `{instance_key}` has invalid restart policy (attempts {max_attempts}, window {window:?})"
),
Self::ActivationCycle { instances } => write!(
formatter,
"required Capability activation cycle: {}",
instances.join(" -> ")
),
}
}
}
impl std::error::Error for PlanResolutionError {}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct ResolvedAppPlan {
schema_version: u32,
module_instances: Vec<ModuleInstancePlan>,
capability_bindings: Vec<CapabilityBinding>,
#[serde(default = "default_execution_lanes")]
execution_lanes: Vec<ExecutionLanePlan>,
}
impl ResolvedAppPlan {
pub fn empty() -> Self {
Self {
schema_version: PLAN_SCHEMA_VERSION,
module_instances: Vec::new(),
capability_bindings: Vec::new(),
execution_lanes: default_execution_lanes(),
}
}
pub fn new(
mut module_instances: Vec<ModuleInstancePlan>,
mut capability_bindings: Vec<CapabilityBinding>,
) -> Self {
sort_module_instances(&mut module_instances);
sort_bindings(&mut capability_bindings);
Self {
schema_version: PLAN_SCHEMA_VERSION,
module_instances,
capability_bindings,
execution_lanes: default_execution_lanes(),
}
}
pub const fn with_schema_version(schema_version: u32) -> Self {
Self {
schema_version,
module_instances: Vec::new(),
capability_bindings: Vec::new(),
execution_lanes: Vec::new(),
}
}
pub fn validate(&self) -> Result<(), PlanResolutionError> {
if self.schema_version != PLAN_SCHEMA_VERSION {
return Err(PlanResolutionError::UnsupportedSchemaVersion {
expected: PLAN_SCHEMA_VERSION,
actual: self.schema_version,
});
}
validate_execution_lanes(&self.execution_lanes, &self.module_instances)?;
resolve_parts(&self.module_instances, &self.capability_bindings).map(|_| ())
}
pub fn activation_order(&self) -> Result<Vec<String>, PlanResolutionError> {
if self.schema_version != PLAN_SCHEMA_VERSION {
return Err(PlanResolutionError::UnsupportedSchemaVersion {
expected: PLAN_SCHEMA_VERSION,
actual: self.schema_version,
});
}
validate_execution_lanes(&self.execution_lanes, &self.module_instances)?;
let (instances, bindings) =
resolve_parts(&self.module_instances, &self.capability_bindings)?;
activation_order_for(&instances, &bindings)
.map_err(|instances| PlanResolutionError::ActivationCycle { instances })
}
pub const fn schema_version(&self) -> u32 {
self.schema_version
}
pub fn module_instances(&self) -> &[ModuleInstancePlan] {
&self.module_instances
}
pub fn capability_bindings(&self) -> &[CapabilityBinding] {
&self.capability_bindings
}
pub fn execution_lanes(&self) -> &[ExecutionLanePlan] {
&self.execution_lanes
}
pub fn request_admission_for(
&self,
binding: &CapabilityBinding,
operation: &str,
) -> RequestAdmissionPlan {
if binding.has_explicit_admission() {
return binding.admission();
}
self.module_instances
.iter()
.find(|instance| instance.instance_key() == binding.provider_instance())
.and_then(|provider| {
provider
.provided_capabilities()
.iter()
.find(|endpoint| endpoint.capability_id() == binding.capability_id())
})
.and_then(|endpoint| endpoint.operation_admission(operation))
.unwrap_or_else(|| binding.admission())
}
pub fn event_admission_for(&self, binding: &CapabilityBinding) -> EventAdmissionPlan {
if binding.has_explicit_event_admission() {
return binding.event_admission();
}
self.module_instances
.iter()
.find(|instance| instance.instance_key() == binding.provider_instance())
.and_then(|provider| {
provider
.provided_capabilities()
.iter()
.find(|endpoint| endpoint.capability_id() == binding.capability_id())
})
.and_then(CapabilityEndpointPlan::event_admission)
.unwrap_or_else(|| binding.event_admission())
}
pub fn module_instance(&self, instance_key: &str) -> Option<&ModuleInstancePlan> {
self.module_instances
.iter()
.find(|instance| instance.instance_key() == instance_key)
}
pub fn restart_policy_for(&self, instance_key: &str) -> Option<RestartPolicy> {
self.module_instance(instance_key)
.map(ModuleInstancePlan::restart_policy)
}
pub fn criticality_for(&self, instance_key: &str) -> Option<ModuleCriticality> {
self.module_instance(instance_key)
.map(ModuleInstancePlan::criticality)
}
pub fn module_instance_is_required(&self, instance_key: &str) -> bool {
self.capability_bindings.iter().any(|binding| {
binding.provider_instance() == instance_key
&& self
.module_instance(binding.consumer_instance())
.is_some_and(|consumer| {
consumer.required_capabilities().iter().any(|requirement| {
requirement.capability_id() == binding.capability_id()
&& requirement.cardinality() == CapabilityCardinality::One
})
})
})
}
}