1pub mod authoring;
4
5use std::collections::BTreeMap;
6
7use serde::{Deserialize, Serialize};
8
9mod contract;
10pub use contract::{CapabilityBinding, CapabilityRequirementPlan, PluginInstancePlan};
11mod error;
12mod execution;
13mod policy;
14mod resolution;
15mod schema;
16pub use schema::TerminalPolicy;
17
18pub use error::PlanResolutionError;
19pub use execution::{ExecutionClassId, ExecutionLaneId, ExecutionLanePlan};
20pub use policy::{
21 CapabilityCardinality, CapabilityOperationKind, EventAdmissionPlan, PluginCriticality,
22 RequestAdmissionPlan, RestartMode, RestartPolicy,
23};
24use resolution::{
25 activation_order_for, resolve_parts, sort_bindings, sort_plugin_instances,
26 sorted_execution_lanes, validate_execution_lanes,
27};
28
29pub const PLAN_SCHEMA_VERSION: u32 = 3;
31
32pub const DEFAULT_REQUEST_QUEUE_CAPACITY: usize = 16;
34
35pub const DEFAULT_REQUEST_MAX_CONCURRENCY: usize = 1;
37
38pub const DEFAULT_EVENT_QUEUE_CAPACITY: usize = 16;
40
41fn default_execution_lanes() -> Vec<ExecutionLanePlan> {
42 vec![ExecutionLanePlan::new("main")]
43}
44
45#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
47pub struct CapabilityEndpointPlan {
48 capability_id: String,
49 descriptor_version: String,
50 operations: Vec<String>,
51 operation_kinds: BTreeMap<String, CapabilityOperationKind>,
52 default_admission: Option<RequestAdmissionPlan>,
53 operation_admissions: BTreeMap<String, RequestAdmissionPlan>,
54 event_admission: Option<EventAdmissionPlan>,
55 #[serde(default)]
56 cross_lane_transfer: bool,
57}
58
59impl CapabilityEndpointPlan {
60 pub fn new(
62 capability_id: impl Into<String>,
63 descriptor_version: impl Into<String>,
64 operations: impl IntoIterator<Item = impl Into<String>>,
65 ) -> Self {
66 Self {
67 capability_id: capability_id.into(),
68 descriptor_version: descriptor_version.into(),
69 operations: operations.into_iter().map(Into::into).collect(),
70 operation_kinds: BTreeMap::new(),
71 default_admission: None,
72 operation_admissions: BTreeMap::new(),
73 event_admission: None,
74 cross_lane_transfer: false,
75 }
76 }
77
78 #[must_use]
80 pub fn with_admission(mut self, admission: RequestAdmissionPlan) -> Self {
81 self.default_admission = Some(admission);
82 self
83 }
84
85 #[must_use]
87 pub fn with_limits(self, queue_capacity: usize, max_concurrency: usize) -> Self {
88 self.with_admission(RequestAdmissionPlan::new(queue_capacity, max_concurrency))
89 }
90
91 #[must_use]
93 pub fn with_operation_kind(
94 mut self,
95 operation: impl Into<String>,
96 kind: CapabilityOperationKind,
97 ) -> Self {
98 self.operation_kinds.insert(operation.into(), kind);
99 self
100 }
101
102 #[must_use]
104 pub fn with_stream_operation(self, operation: impl Into<String>) -> Self {
105 self.with_operation_kind(operation, CapabilityOperationKind::Stream)
106 }
107
108 #[must_use]
110 pub fn with_event_operation(self, operation: impl Into<String>) -> Self {
111 self.with_operation_kind(operation, CapabilityOperationKind::Event)
112 }
113
114 #[must_use]
116 pub fn with_event_admission(mut self, admission: EventAdmissionPlan) -> Self {
117 self.event_admission = Some(admission);
118 self
119 }
120
121 #[must_use]
123 pub fn with_event_capacity(self, capacity: usize) -> Self {
124 self.with_event_admission(EventAdmissionPlan::new(capacity))
125 }
126
127 #[must_use]
129 pub const fn with_cross_lane_transfer(mut self) -> Self {
130 self.cross_lane_transfer = true;
131 self
132 }
133
134 #[must_use]
136 pub fn with_operation_admission(
137 mut self,
138 operation: impl Into<String>,
139 admission: RequestAdmissionPlan,
140 ) -> Self {
141 self.operation_admissions
142 .insert(operation.into(), admission);
143 self
144 }
145
146 #[must_use]
148 pub fn with_operation_limits(
149 self,
150 operation: impl Into<String>,
151 queue_capacity: usize,
152 max_concurrency: usize,
153 ) -> Self {
154 self.with_operation_admission(
155 operation,
156 RequestAdmissionPlan::new(queue_capacity, max_concurrency),
157 )
158 }
159
160 pub fn capability_id(&self) -> &str {
162 &self.capability_id
163 }
164
165 pub fn descriptor_version(&self) -> &str {
167 &self.descriptor_version
168 }
169
170 pub fn operations(&self) -> &[String] {
172 &self.operations
173 }
174
175 pub fn operation_kind(&self, operation: &str) -> Option<CapabilityOperationKind> {
177 self.operations
178 .iter()
179 .any(|declared| declared == operation)
180 .then(|| {
181 self.operation_kinds
182 .get(operation)
183 .copied()
184 .unwrap_or(CapabilityOperationKind::Request)
185 })
186 }
187
188 pub fn stream_operations(&self) -> Vec<&str> {
190 self.operations
191 .iter()
192 .filter(|operation| {
193 self.operation_kind(operation) == Some(CapabilityOperationKind::Stream)
194 })
195 .map(String::as_str)
196 .collect()
197 }
198
199 pub fn request_operations(&self) -> Vec<&str> {
201 self.operations
202 .iter()
203 .filter(|operation| {
204 self.operation_kind(operation) == Some(CapabilityOperationKind::Request)
205 })
206 .map(String::as_str)
207 .collect()
208 }
209
210 pub fn event_operations(&self) -> Vec<&str> {
212 self.operations
213 .iter()
214 .filter(|operation| {
215 self.operation_kind(operation) == Some(CapabilityOperationKind::Event)
216 })
217 .map(String::as_str)
218 .collect()
219 }
220
221 pub const fn event_admission(&self) -> Option<EventAdmissionPlan> {
223 self.event_admission
224 }
225
226 pub const fn supports_cross_lane_transfer(&self) -> bool {
228 self.cross_lane_transfer
229 }
230
231 pub fn default_admission(&self) -> Option<RequestAdmissionPlan> {
233 self.default_admission
234 }
235
236 pub fn operation_admissions(&self) -> &BTreeMap<String, RequestAdmissionPlan> {
238 &self.operation_admissions
239 }
240
241 pub fn operation_admission(&self, operation: &str) -> Option<RequestAdmissionPlan> {
243 self.operation_admissions
244 .get(operation)
245 .copied()
246 .or(self.default_admission)
247 }
248}
249
250#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
252pub struct AppComposition {
253 plugin_instances: Vec<PluginInstancePlan>,
254 capability_bindings: Vec<CapabilityBinding>,
255 #[serde(default = "default_execution_lanes")]
256 execution_lanes: Vec<ExecutionLanePlan>,
257}
258
259impl AppComposition {
260 pub fn new(
262 plugin_instances: Vec<PluginInstancePlan>,
263 capability_bindings: Vec<CapabilityBinding>,
264 ) -> Self {
265 Self {
266 plugin_instances,
267 capability_bindings,
268 execution_lanes: default_execution_lanes(),
269 }
270 }
271
272 #[must_use]
274 pub fn with_execution_lanes(mut self, execution_lanes: Vec<ExecutionLanePlan>) -> Self {
275 self.execution_lanes = execution_lanes;
276 self
277 }
278
279 pub fn resolve(&self) -> Result<ResolvedAppPlan, PlanResolutionError> {
281 validate_execution_lanes(&self.execution_lanes, &self.plugin_instances)?;
282 resolve_parts(&self.plugin_instances, &self.capability_bindings).map(
283 |(plugin_instances, capability_bindings)| ResolvedAppPlan {
284 terminal_policy: TerminalPolicy::RequiredPath,
285 schema_version: PLAN_SCHEMA_VERSION,
286 plugin_instances,
287 capability_bindings,
288 execution_lanes: sorted_execution_lanes(&self.execution_lanes),
289 },
290 )
291 }
292
293 pub fn plugin_instances(&self) -> &[PluginInstancePlan] {
295 &self.plugin_instances
296 }
297
298 pub fn capability_bindings(&self) -> &[CapabilityBinding] {
300 &self.capability_bindings
301 }
302
303 pub fn execution_lanes(&self) -> &[ExecutionLanePlan] {
305 &self.execution_lanes
306 }
307}
308
309#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
311#[serde(try_from = "schema::PlanWire")]
312pub struct ResolvedAppPlan {
313 terminal_policy: TerminalPolicy,
314 schema_version: u32,
315 plugin_instances: Vec<PluginInstancePlan>,
316 capability_bindings: Vec<CapabilityBinding>,
317 #[serde(default = "default_execution_lanes")]
318 execution_lanes: Vec<ExecutionLanePlan>,
319}
320
321impl ResolvedAppPlan {
322 pub fn empty() -> Self {
324 Self {
325 terminal_policy: TerminalPolicy::RequiredPath,
326 schema_version: PLAN_SCHEMA_VERSION,
327 plugin_instances: Vec::new(),
328 capability_bindings: Vec::new(),
329 execution_lanes: default_execution_lanes(),
330 }
331 }
332
333 pub fn new(
335 mut plugin_instances: Vec<PluginInstancePlan>,
336 mut capability_bindings: Vec<CapabilityBinding>,
337 ) -> Self {
338 sort_plugin_instances(&mut plugin_instances);
339 sort_bindings(&mut capability_bindings);
340 Self {
341 terminal_policy: TerminalPolicy::RequiredPath,
342 schema_version: PLAN_SCHEMA_VERSION,
343 plugin_instances,
344 capability_bindings,
345 execution_lanes: default_execution_lanes(),
346 }
347 }
348
349 pub const fn with_schema_version(schema_version: u32) -> Self {
353 Self {
354 terminal_policy: TerminalPolicy::RequiredPath,
355 schema_version,
356 plugin_instances: Vec::new(),
357 capability_bindings: Vec::new(),
358 execution_lanes: Vec::new(),
359 }
360 }
361
362 pub fn validate(&self) -> Result<(), PlanResolutionError> {
364 if self.schema_version != PLAN_SCHEMA_VERSION {
365 return Err(PlanResolutionError::UnsupportedSchemaVersion {
366 expected: PLAN_SCHEMA_VERSION,
367 actual: self.schema_version,
368 });
369 }
370 validate_execution_lanes(&self.execution_lanes, &self.plugin_instances)?;
371 let (instances, bindings) =
372 resolve_parts(&self.plugin_instances, &self.capability_bindings)?;
373 self.terminal_policy.validate(&instances, &bindings)
374 }
375
376 pub fn activation_order(&self) -> Result<Vec<String>, PlanResolutionError> {
381 if self.schema_version != PLAN_SCHEMA_VERSION {
382 return Err(PlanResolutionError::UnsupportedSchemaVersion {
383 expected: PLAN_SCHEMA_VERSION,
384 actual: self.schema_version,
385 });
386 }
387 validate_execution_lanes(&self.execution_lanes, &self.plugin_instances)?;
388 let (instances, bindings) =
389 resolve_parts(&self.plugin_instances, &self.capability_bindings)?;
390 self.terminal_policy.validate(&instances, &bindings)?;
391 activation_order_for(&instances, &bindings)
392 .map_err(|instances| PlanResolutionError::ActivationCycle { instances })
393 }
394
395 pub const fn terminal_policy(&self) -> &TerminalPolicy {
397 &self.terminal_policy
398 }
399
400 #[must_use]
402 pub fn with_terminal_policy(mut self, policy: TerminalPolicy) -> Self {
403 self.terminal_policy = policy;
404 self
405 }
406
407 pub const fn schema_version(&self) -> u32 {
409 self.schema_version
410 }
411
412 pub fn plugin_instances(&self) -> &[PluginInstancePlan] {
414 &self.plugin_instances
415 }
416
417 pub fn capability_bindings(&self) -> &[CapabilityBinding] {
419 &self.capability_bindings
420 }
421
422 pub fn execution_lanes(&self) -> &[ExecutionLanePlan] {
424 &self.execution_lanes
425 }
426
427 pub fn request_admission_for(
429 &self,
430 binding: &CapabilityBinding,
431 operation: &str,
432 ) -> RequestAdmissionPlan {
433 if binding.has_explicit_admission() {
434 return binding.admission();
435 }
436
437 self.plugin_instances
438 .iter()
439 .find(|instance| instance.instance_key() == binding.provider_instance())
440 .and_then(|provider| {
441 provider
442 .provided_capabilities()
443 .iter()
444 .find(|endpoint| endpoint.capability_id() == binding.capability_id())
445 })
446 .and_then(|endpoint| endpoint.operation_admission(operation))
447 .unwrap_or_else(|| binding.admission())
448 }
449
450 pub fn event_admission_for(&self, binding: &CapabilityBinding) -> EventAdmissionPlan {
452 if binding.has_explicit_event_admission() {
453 return binding.event_admission();
454 }
455
456 self.plugin_instances
457 .iter()
458 .find(|instance| instance.instance_key() == binding.provider_instance())
459 .and_then(|provider| {
460 provider
461 .provided_capabilities()
462 .iter()
463 .find(|endpoint| endpoint.capability_id() == binding.capability_id())
464 })
465 .and_then(CapabilityEndpointPlan::event_admission)
466 .unwrap_or_else(|| binding.event_admission())
467 }
468
469 pub fn plugin_instance(&self, instance_key: &str) -> Option<&PluginInstancePlan> {
471 self.plugin_instances
472 .iter()
473 .find(|instance| instance.instance_key() == instance_key)
474 }
475
476 pub fn restart_policy_for(&self, instance_key: &str) -> Option<RestartPolicy> {
478 self.plugin_instance(instance_key)
479 .map(PluginInstancePlan::restart_policy)
480 }
481
482 pub fn criticality_for(&self, instance_key: &str) -> Option<PluginCriticality> {
484 self.plugin_instance(instance_key)
485 .map(PluginInstancePlan::criticality)
486 }
487
488 pub fn plugin_instance_is_required(&self, instance_key: &str) -> bool {
490 self.capability_bindings.iter().any(|binding| {
491 binding.provider_instance() == instance_key
492 && self
493 .plugin_instance(binding.consumer_instance())
494 .is_some_and(|consumer| {
495 consumer.required_capabilities().iter().any(|requirement| {
496 requirement.requirement_id() == binding.requirement_id()
497 && requirement.cardinality() == CapabilityCardinality::One
498 })
499 })
500 })
501 }
502
503 pub fn plugin_instance_is_terminal(&self, instance_key: &str) -> bool {
506 match &self.terminal_policy {
507 TerminalPolicy::RequiredPath => self.plugin_instance_is_required(instance_key),
508 TerminalPolicy::HostEssential { closure, .. } => closure
509 .binary_search_by(|candidate| candidate.as_str().cmp(instance_key))
510 .is_ok(),
511 }
512 }
513}