pub mod authoring;
use std::collections::BTreeMap;
use serde::{Deserialize, Serialize};
mod contract;
pub use contract::{CapabilityBinding, CapabilityRequirementPlan, PluginInstancePlan};
mod error;
mod execution;
mod policy;
mod resolution;
mod schema;
pub use schema::TerminalPolicy;
pub use error::PlanResolutionError;
pub use execution::{ExecutionClassId, ExecutionLaneId, ExecutionLanePlan};
pub use policy::{
CapabilityCardinality, CapabilityOperationKind, EventAdmissionPlan, PluginCriticality,
RequestAdmissionPlan, RestartMode, RestartPolicy,
};
use resolution::{
activation_order_for, resolve_parts, sort_bindings, sort_plugin_instances,
sorted_execution_lanes, validate_execution_lanes,
};
pub const PLAN_SCHEMA_VERSION: u32 = 3;
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, 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 AppComposition {
plugin_instances: Vec<PluginInstancePlan>,
capability_bindings: Vec<CapabilityBinding>,
#[serde(default = "default_execution_lanes")]
execution_lanes: Vec<ExecutionLanePlan>,
}
impl AppComposition {
pub fn new(
plugin_instances: Vec<PluginInstancePlan>,
capability_bindings: Vec<CapabilityBinding>,
) -> Self {
Self {
plugin_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.plugin_instances)?;
resolve_parts(&self.plugin_instances, &self.capability_bindings).map(
|(plugin_instances, capability_bindings)| ResolvedAppPlan {
terminal_policy: TerminalPolicy::RequiredPath,
schema_version: PLAN_SCHEMA_VERSION,
plugin_instances,
capability_bindings,
execution_lanes: sorted_execution_lanes(&self.execution_lanes),
},
)
}
pub fn plugin_instances(&self) -> &[PluginInstancePlan] {
&self.plugin_instances
}
pub fn capability_bindings(&self) -> &[CapabilityBinding] {
&self.capability_bindings
}
pub fn execution_lanes(&self) -> &[ExecutionLanePlan] {
&self.execution_lanes
}
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(try_from = "schema::PlanWire")]
pub struct ResolvedAppPlan {
terminal_policy: TerminalPolicy,
schema_version: u32,
plugin_instances: Vec<PluginInstancePlan>,
capability_bindings: Vec<CapabilityBinding>,
#[serde(default = "default_execution_lanes")]
execution_lanes: Vec<ExecutionLanePlan>,
}
impl ResolvedAppPlan {
pub fn empty() -> Self {
Self {
terminal_policy: TerminalPolicy::RequiredPath,
schema_version: PLAN_SCHEMA_VERSION,
plugin_instances: Vec::new(),
capability_bindings: Vec::new(),
execution_lanes: default_execution_lanes(),
}
}
pub fn new(
mut plugin_instances: Vec<PluginInstancePlan>,
mut capability_bindings: Vec<CapabilityBinding>,
) -> Self {
sort_plugin_instances(&mut plugin_instances);
sort_bindings(&mut capability_bindings);
Self {
terminal_policy: TerminalPolicy::RequiredPath,
schema_version: PLAN_SCHEMA_VERSION,
plugin_instances,
capability_bindings,
execution_lanes: default_execution_lanes(),
}
}
pub const fn with_schema_version(schema_version: u32) -> Self {
Self {
terminal_policy: TerminalPolicy::RequiredPath,
schema_version,
plugin_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.plugin_instances)?;
let (instances, bindings) =
resolve_parts(&self.plugin_instances, &self.capability_bindings)?;
self.terminal_policy.validate(&instances, &bindings)
}
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.plugin_instances)?;
let (instances, bindings) =
resolve_parts(&self.plugin_instances, &self.capability_bindings)?;
self.terminal_policy.validate(&instances, &bindings)?;
activation_order_for(&instances, &bindings)
.map_err(|instances| PlanResolutionError::ActivationCycle { instances })
}
pub const fn terminal_policy(&self) -> &TerminalPolicy {
&self.terminal_policy
}
#[must_use]
pub fn with_terminal_policy(mut self, policy: TerminalPolicy) -> Self {
self.terminal_policy = policy;
self
}
pub const fn schema_version(&self) -> u32 {
self.schema_version
}
pub fn plugin_instances(&self) -> &[PluginInstancePlan] {
&self.plugin_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.plugin_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.plugin_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 plugin_instance(&self, instance_key: &str) -> Option<&PluginInstancePlan> {
self.plugin_instances
.iter()
.find(|instance| instance.instance_key() == instance_key)
}
pub fn restart_policy_for(&self, instance_key: &str) -> Option<RestartPolicy> {
self.plugin_instance(instance_key)
.map(PluginInstancePlan::restart_policy)
}
pub fn criticality_for(&self, instance_key: &str) -> Option<PluginCriticality> {
self.plugin_instance(instance_key)
.map(PluginInstancePlan::criticality)
}
pub fn plugin_instance_is_required(&self, instance_key: &str) -> bool {
self.capability_bindings.iter().any(|binding| {
binding.provider_instance() == instance_key
&& self
.plugin_instance(binding.consumer_instance())
.is_some_and(|consumer| {
consumer.required_capabilities().iter().any(|requirement| {
requirement.requirement_id() == binding.requirement_id()
&& requirement.cardinality() == CapabilityCardinality::One
})
})
})
}
pub fn plugin_instance_is_terminal(&self, instance_key: &str) -> bool {
match &self.terminal_policy {
TerminalPolicy::RequiredPath => self.plugin_instance_is_required(instance_key),
TerminalPolicy::HostEssential { closure, .. } => closure
.binary_search_by(|candidate| candidate.as_str().cmp(instance_key))
.is_ok(),
}
}
}