1pub mod authoring;
4
5use std::collections::BTreeMap;
6
7use serde::{Deserialize, Serialize};
8
9mod error;
10mod execution;
11mod policy;
12mod resolution;
13
14pub use error::PlanResolutionError;
15pub use execution::{ExecutionClassId, ExecutionLaneId, ExecutionLanePlan};
16pub use policy::{
17 CapabilityCardinality, CapabilityOperationKind, EventAdmissionPlan, ModuleCriticality,
18 RequestAdmissionPlan, RestartMode, RestartPolicy,
19};
20use resolution::{
21 activation_order_for, resolve_parts, sort_bindings, sort_module_instances,
22 sorted_execution_lanes, validate_execution_lanes,
23};
24
25pub const PLAN_SCHEMA_VERSION: u32 = 1;
27
28pub const DEFAULT_REQUEST_QUEUE_CAPACITY: usize = 16;
30
31pub const DEFAULT_REQUEST_MAX_CONCURRENCY: usize = 1;
33
34pub const DEFAULT_EVENT_QUEUE_CAPACITY: usize = 16;
36
37fn default_execution_lanes() -> Vec<ExecutionLanePlan> {
38 vec![ExecutionLanePlan::new("main")]
39}
40
41#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
43pub struct CapabilityRequirementPlan {
44 capability_id: String,
45 descriptor_version: String,
46 cardinality: CapabilityCardinality,
47}
48
49impl CapabilityRequirementPlan {
50 pub fn new(
52 capability_id: impl Into<String>,
53 descriptor_version: impl Into<String>,
54 cardinality: CapabilityCardinality,
55 ) -> Self {
56 Self {
57 capability_id: capability_id.into(),
58 descriptor_version: descriptor_version.into(),
59 cardinality,
60 }
61 }
62
63 pub fn one(capability_id: impl Into<String>, descriptor_version: impl Into<String>) -> Self {
65 Self::new(
66 capability_id,
67 descriptor_version,
68 CapabilityCardinality::One,
69 )
70 }
71
72 pub fn optional(
74 capability_id: impl Into<String>,
75 descriptor_version: impl Into<String>,
76 ) -> Self {
77 Self::new(
78 capability_id,
79 descriptor_version,
80 CapabilityCardinality::Optional,
81 )
82 }
83
84 pub fn many(capability_id: impl Into<String>, descriptor_version: impl Into<String>) -> Self {
86 Self::new(
87 capability_id,
88 descriptor_version,
89 CapabilityCardinality::Many,
90 )
91 }
92
93 pub fn capability_id(&self) -> &str {
95 &self.capability_id
96 }
97
98 pub fn descriptor_version(&self) -> &str {
100 &self.descriptor_version
101 }
102
103 pub const fn cardinality(&self) -> CapabilityCardinality {
105 self.cardinality
106 }
107}
108
109#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
111pub struct CapabilityEndpointPlan {
112 capability_id: String,
113 descriptor_version: String,
114 operations: Vec<String>,
115 operation_kinds: BTreeMap<String, CapabilityOperationKind>,
116 default_admission: Option<RequestAdmissionPlan>,
117 operation_admissions: BTreeMap<String, RequestAdmissionPlan>,
118 event_admission: Option<EventAdmissionPlan>,
119 #[serde(default)]
120 cross_lane_transfer: bool,
121}
122
123impl CapabilityEndpointPlan {
124 pub fn new(
126 capability_id: impl Into<String>,
127 descriptor_version: impl Into<String>,
128 operations: impl IntoIterator<Item = impl Into<String>>,
129 ) -> Self {
130 Self {
131 capability_id: capability_id.into(),
132 descriptor_version: descriptor_version.into(),
133 operations: operations.into_iter().map(Into::into).collect(),
134 operation_kinds: BTreeMap::new(),
135 default_admission: None,
136 operation_admissions: BTreeMap::new(),
137 event_admission: None,
138 cross_lane_transfer: false,
139 }
140 }
141
142 #[must_use]
144 pub fn with_admission(mut self, admission: RequestAdmissionPlan) -> Self {
145 self.default_admission = Some(admission);
146 self
147 }
148
149 #[must_use]
151 pub fn with_limits(self, queue_capacity: usize, max_concurrency: usize) -> Self {
152 self.with_admission(RequestAdmissionPlan::new(queue_capacity, max_concurrency))
153 }
154
155 #[must_use]
157 pub fn with_operation_kind(
158 mut self,
159 operation: impl Into<String>,
160 kind: CapabilityOperationKind,
161 ) -> Self {
162 self.operation_kinds.insert(operation.into(), kind);
163 self
164 }
165
166 #[must_use]
168 pub fn with_stream_operation(self, operation: impl Into<String>) -> Self {
169 self.with_operation_kind(operation, CapabilityOperationKind::Stream)
170 }
171
172 #[must_use]
174 pub fn with_event_operation(self, operation: impl Into<String>) -> Self {
175 self.with_operation_kind(operation, CapabilityOperationKind::Event)
176 }
177
178 #[must_use]
180 pub fn with_event_admission(mut self, admission: EventAdmissionPlan) -> Self {
181 self.event_admission = Some(admission);
182 self
183 }
184
185 #[must_use]
187 pub fn with_event_capacity(self, capacity: usize) -> Self {
188 self.with_event_admission(EventAdmissionPlan::new(capacity))
189 }
190
191 #[must_use]
193 pub const fn with_cross_lane_transfer(mut self) -> Self {
194 self.cross_lane_transfer = true;
195 self
196 }
197
198 #[must_use]
200 pub fn with_operation_admission(
201 mut self,
202 operation: impl Into<String>,
203 admission: RequestAdmissionPlan,
204 ) -> Self {
205 self.operation_admissions
206 .insert(operation.into(), admission);
207 self
208 }
209
210 #[must_use]
212 pub fn with_operation_limits(
213 self,
214 operation: impl Into<String>,
215 queue_capacity: usize,
216 max_concurrency: usize,
217 ) -> Self {
218 self.with_operation_admission(
219 operation,
220 RequestAdmissionPlan::new(queue_capacity, max_concurrency),
221 )
222 }
223
224 pub fn capability_id(&self) -> &str {
226 &self.capability_id
227 }
228
229 pub fn descriptor_version(&self) -> &str {
231 &self.descriptor_version
232 }
233
234 pub fn operations(&self) -> &[String] {
236 &self.operations
237 }
238
239 pub fn operation_kind(&self, operation: &str) -> Option<CapabilityOperationKind> {
241 self.operations
242 .iter()
243 .any(|declared| declared == operation)
244 .then(|| {
245 self.operation_kinds
246 .get(operation)
247 .copied()
248 .unwrap_or(CapabilityOperationKind::Request)
249 })
250 }
251
252 pub fn stream_operations(&self) -> Vec<&str> {
254 self.operations
255 .iter()
256 .filter(|operation| {
257 self.operation_kind(operation) == Some(CapabilityOperationKind::Stream)
258 })
259 .map(String::as_str)
260 .collect()
261 }
262
263 pub fn request_operations(&self) -> Vec<&str> {
265 self.operations
266 .iter()
267 .filter(|operation| {
268 self.operation_kind(operation) == Some(CapabilityOperationKind::Request)
269 })
270 .map(String::as_str)
271 .collect()
272 }
273
274 pub fn event_operations(&self) -> Vec<&str> {
276 self.operations
277 .iter()
278 .filter(|operation| {
279 self.operation_kind(operation) == Some(CapabilityOperationKind::Event)
280 })
281 .map(String::as_str)
282 .collect()
283 }
284
285 pub const fn event_admission(&self) -> Option<EventAdmissionPlan> {
287 self.event_admission
288 }
289
290 pub const fn supports_cross_lane_transfer(&self) -> bool {
292 self.cross_lane_transfer
293 }
294
295 pub fn default_admission(&self) -> Option<RequestAdmissionPlan> {
297 self.default_admission
298 }
299
300 pub fn operation_admissions(&self) -> &BTreeMap<String, RequestAdmissionPlan> {
302 &self.operation_admissions
303 }
304
305 pub fn operation_admission(&self, operation: &str) -> Option<RequestAdmissionPlan> {
307 self.operation_admissions
308 .get(operation)
309 .copied()
310 .or(self.default_admission)
311 }
312}
313
314#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
316pub struct ModuleInstancePlan {
317 instance_key: String,
318 package_id: String,
319 entrypoint: String,
320 configuration: String,
321 provided_capabilities: Vec<CapabilityEndpointPlan>,
322 required_capabilities: Vec<CapabilityRequirementPlan>,
323 execution_class: ExecutionClassId,
324 package_revision: String,
325 restart_policy: RestartPolicy,
326 criticality: ModuleCriticality,
327 #[serde(default)]
328 execution_lane: ExecutionLaneId,
329}
330
331impl ModuleInstancePlan {
332 pub fn new(instance_key: impl Into<String>, package_id: impl Into<String>) -> Self {
334 Self {
335 instance_key: instance_key.into(),
336 package_id: package_id.into(),
337 entrypoint: "default".to_owned(),
338 configuration: "{}".to_owned(),
339 provided_capabilities: Vec::new(),
340 required_capabilities: Vec::new(),
341 execution_class: ExecutionClassId::native_rust(),
342 package_revision: String::new(),
343 restart_policy: RestartPolicy::default(),
344 criticality: ModuleCriticality::default(),
345 execution_lane: ExecutionLaneId::default(),
346 }
347 }
348
349 #[must_use]
351 pub fn with_entrypoint(mut self, entrypoint: impl Into<String>) -> Self {
352 self.entrypoint = entrypoint.into();
353 self
354 }
355
356 #[must_use]
358 pub fn with_configuration(mut self, configuration: impl Into<String>) -> Self {
359 self.configuration = configuration.into();
360 self
361 }
362
363 #[must_use]
365 pub fn with_capability(mut self, capability: CapabilityEndpointPlan) -> Self {
366 self.provided_capabilities.push(capability);
367 self
368 }
369
370 #[must_use]
372 pub fn with_requirement(mut self, requirement: CapabilityRequirementPlan) -> Self {
373 self.required_capabilities.push(requirement);
374 self
375 }
376
377 #[must_use]
379 pub fn with_required_capability(self, requirement: CapabilityRequirementPlan) -> Self {
380 self.with_requirement(requirement)
381 }
382
383 #[must_use]
385 pub fn with_execution_class(mut self, execution_class: ExecutionClassId) -> Self {
386 self.execution_class = execution_class;
387 self
388 }
389
390 #[must_use]
392 pub fn with_execution_lane(mut self, execution_lane: ExecutionLaneId) -> Self {
393 self.execution_lane = execution_lane;
394 self
395 }
396
397 #[must_use]
399 pub fn with_package_revision(mut self, revision: impl Into<String>) -> Self {
400 self.package_revision = revision.into();
401 self
402 }
403
404 #[must_use]
406 pub fn with_restart_policy(mut self, restart_policy: RestartPolicy) -> Self {
407 self.restart_policy = restart_policy;
408 self
409 }
410
411 #[must_use]
413 pub fn with_criticality(mut self, criticality: ModuleCriticality) -> Self {
414 self.criticality = criticality;
415 self
416 }
417
418 pub fn instance_key(&self) -> &str {
420 &self.instance_key
421 }
422
423 pub fn package_id(&self) -> &str {
425 &self.package_id
426 }
427
428 pub fn entrypoint(&self) -> &str {
430 &self.entrypoint
431 }
432
433 pub fn configuration(&self) -> &str {
435 &self.configuration
436 }
437
438 pub fn provided_capabilities(&self) -> &[CapabilityEndpointPlan] {
440 &self.provided_capabilities
441 }
442
443 pub fn required_capabilities(&self) -> &[CapabilityRequirementPlan] {
445 &self.required_capabilities
446 }
447
448 pub fn execution_class(&self) -> &ExecutionClassId {
450 &self.execution_class
451 }
452
453 pub const fn execution_lane(&self) -> &ExecutionLaneId {
455 &self.execution_lane
456 }
457
458 pub fn package_revision(&self) -> &str {
460 &self.package_revision
461 }
462
463 pub const fn restart_policy(&self) -> RestartPolicy {
465 self.restart_policy
466 }
467
468 pub const fn criticality(&self) -> ModuleCriticality {
470 self.criticality
471 }
472}
473
474#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
476pub struct CapabilityBinding {
477 consumer_instance: String,
478 capability_id: String,
479 descriptor_version: String,
480 provider_instance: String,
481 provider_order: usize,
482 admission: RequestAdmissionPlan,
483 admission_explicit: bool,
484 event_admission: EventAdmissionPlan,
485 event_admission_explicit: bool,
486}
487
488impl CapabilityBinding {
489 pub fn new(
491 consumer_instance: impl Into<String>,
492 capability_id: impl Into<String>,
493 descriptor_version: impl Into<String>,
494 provider_instance: impl Into<String>,
495 ) -> Self {
496 Self {
497 consumer_instance: consumer_instance.into(),
498 capability_id: capability_id.into(),
499 descriptor_version: descriptor_version.into(),
500 provider_instance: provider_instance.into(),
501 provider_order: 0,
502 admission: RequestAdmissionPlan::default(),
503 admission_explicit: false,
504 event_admission: EventAdmissionPlan::default(),
505 event_admission_explicit: false,
506 }
507 }
508
509 #[must_use]
511 pub fn with_admission(mut self, admission: RequestAdmissionPlan) -> Self {
512 self.admission = admission;
513 self.admission_explicit = true;
514 self
515 }
516
517 #[must_use]
519 pub fn with_limits(self, queue_capacity: usize, max_concurrency: usize) -> Self {
520 self.with_admission(RequestAdmissionPlan::new(queue_capacity, max_concurrency))
521 }
522
523 #[must_use]
525 pub fn with_event_admission(mut self, admission: EventAdmissionPlan) -> Self {
526 self.event_admission = admission;
527 self.event_admission_explicit = true;
528 self
529 }
530
531 #[must_use]
533 pub fn with_event_capacity(self, capacity: usize) -> Self {
534 self.with_event_admission(EventAdmissionPlan::new(capacity))
535 }
536
537 fn with_provider_order(mut self, provider_order: usize) -> Self {
538 self.provider_order = provider_order;
539 self
540 }
541
542 pub fn consumer_instance(&self) -> &str {
544 &self.consumer_instance
545 }
546
547 pub fn capability_id(&self) -> &str {
549 &self.capability_id
550 }
551
552 pub fn descriptor_version(&self) -> &str {
554 &self.descriptor_version
555 }
556
557 pub fn provider_instance(&self) -> &str {
559 &self.provider_instance
560 }
561
562 pub const fn provider_order(&self) -> usize {
564 self.provider_order
565 }
566
567 pub const fn admission(&self) -> RequestAdmissionPlan {
569 self.admission
570 }
571
572 pub const fn has_explicit_admission(&self) -> bool {
574 self.admission_explicit
575 }
576
577 pub const fn event_admission(&self) -> EventAdmissionPlan {
579 self.event_admission
580 }
581
582 pub const fn has_explicit_event_admission(&self) -> bool {
584 self.event_admission_explicit
585 }
586}
587
588#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
590pub struct AppComposition {
591 module_instances: Vec<ModuleInstancePlan>,
592 capability_bindings: Vec<CapabilityBinding>,
593 #[serde(default = "default_execution_lanes")]
594 execution_lanes: Vec<ExecutionLanePlan>,
595}
596
597impl AppComposition {
598 pub fn new(
600 module_instances: Vec<ModuleInstancePlan>,
601 capability_bindings: Vec<CapabilityBinding>,
602 ) -> Self {
603 Self {
604 module_instances,
605 capability_bindings,
606 execution_lanes: default_execution_lanes(),
607 }
608 }
609
610 #[must_use]
612 pub fn with_execution_lanes(mut self, execution_lanes: Vec<ExecutionLanePlan>) -> Self {
613 self.execution_lanes = execution_lanes;
614 self
615 }
616
617 pub fn resolve(&self) -> Result<ResolvedAppPlan, PlanResolutionError> {
619 validate_execution_lanes(&self.execution_lanes, &self.module_instances)?;
620 resolve_parts(&self.module_instances, &self.capability_bindings).map(
621 |(module_instances, capability_bindings)| ResolvedAppPlan {
622 schema_version: PLAN_SCHEMA_VERSION,
623 module_instances,
624 capability_bindings,
625 execution_lanes: sorted_execution_lanes(&self.execution_lanes),
626 },
627 )
628 }
629
630 pub fn module_instances(&self) -> &[ModuleInstancePlan] {
632 &self.module_instances
633 }
634
635 pub fn capability_bindings(&self) -> &[CapabilityBinding] {
637 &self.capability_bindings
638 }
639
640 pub fn execution_lanes(&self) -> &[ExecutionLanePlan] {
642 &self.execution_lanes
643 }
644}
645
646#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
648pub struct ResolvedAppPlan {
649 schema_version: u32,
650 module_instances: Vec<ModuleInstancePlan>,
651 capability_bindings: Vec<CapabilityBinding>,
652 #[serde(default = "default_execution_lanes")]
653 execution_lanes: Vec<ExecutionLanePlan>,
654}
655
656impl ResolvedAppPlan {
657 pub fn empty() -> Self {
659 Self {
660 schema_version: PLAN_SCHEMA_VERSION,
661 module_instances: Vec::new(),
662 capability_bindings: Vec::new(),
663 execution_lanes: default_execution_lanes(),
664 }
665 }
666
667 pub fn new(
669 mut module_instances: Vec<ModuleInstancePlan>,
670 mut capability_bindings: Vec<CapabilityBinding>,
671 ) -> Self {
672 sort_module_instances(&mut module_instances);
673 sort_bindings(&mut capability_bindings);
674 Self {
675 schema_version: PLAN_SCHEMA_VERSION,
676 module_instances,
677 capability_bindings,
678 execution_lanes: default_execution_lanes(),
679 }
680 }
681
682 pub const fn with_schema_version(schema_version: u32) -> Self {
686 Self {
687 schema_version,
688 module_instances: Vec::new(),
689 capability_bindings: Vec::new(),
690 execution_lanes: Vec::new(),
691 }
692 }
693
694 pub fn validate(&self) -> Result<(), PlanResolutionError> {
696 if self.schema_version != PLAN_SCHEMA_VERSION {
697 return Err(PlanResolutionError::UnsupportedSchemaVersion {
698 expected: PLAN_SCHEMA_VERSION,
699 actual: self.schema_version,
700 });
701 }
702 validate_execution_lanes(&self.execution_lanes, &self.module_instances)?;
703 resolve_parts(&self.module_instances, &self.capability_bindings).map(|_| ())
704 }
705
706 pub fn activation_order(&self) -> Result<Vec<String>, PlanResolutionError> {
711 if self.schema_version != PLAN_SCHEMA_VERSION {
712 return Err(PlanResolutionError::UnsupportedSchemaVersion {
713 expected: PLAN_SCHEMA_VERSION,
714 actual: self.schema_version,
715 });
716 }
717 validate_execution_lanes(&self.execution_lanes, &self.module_instances)?;
718 let (instances, bindings) =
719 resolve_parts(&self.module_instances, &self.capability_bindings)?;
720 activation_order_for(&instances, &bindings)
721 .map_err(|instances| PlanResolutionError::ActivationCycle { instances })
722 }
723
724 pub const fn schema_version(&self) -> u32 {
726 self.schema_version
727 }
728
729 pub fn module_instances(&self) -> &[ModuleInstancePlan] {
731 &self.module_instances
732 }
733
734 pub fn capability_bindings(&self) -> &[CapabilityBinding] {
736 &self.capability_bindings
737 }
738
739 pub fn execution_lanes(&self) -> &[ExecutionLanePlan] {
741 &self.execution_lanes
742 }
743
744 pub fn request_admission_for(
746 &self,
747 binding: &CapabilityBinding,
748 operation: &str,
749 ) -> RequestAdmissionPlan {
750 if binding.has_explicit_admission() {
751 return binding.admission();
752 }
753
754 self.module_instances
755 .iter()
756 .find(|instance| instance.instance_key() == binding.provider_instance())
757 .and_then(|provider| {
758 provider
759 .provided_capabilities()
760 .iter()
761 .find(|endpoint| endpoint.capability_id() == binding.capability_id())
762 })
763 .and_then(|endpoint| endpoint.operation_admission(operation))
764 .unwrap_or_else(|| binding.admission())
765 }
766
767 pub fn event_admission_for(&self, binding: &CapabilityBinding) -> EventAdmissionPlan {
769 if binding.has_explicit_event_admission() {
770 return binding.event_admission();
771 }
772
773 self.module_instances
774 .iter()
775 .find(|instance| instance.instance_key() == binding.provider_instance())
776 .and_then(|provider| {
777 provider
778 .provided_capabilities()
779 .iter()
780 .find(|endpoint| endpoint.capability_id() == binding.capability_id())
781 })
782 .and_then(CapabilityEndpointPlan::event_admission)
783 .unwrap_or_else(|| binding.event_admission())
784 }
785
786 pub fn module_instance(&self, instance_key: &str) -> Option<&ModuleInstancePlan> {
788 self.module_instances
789 .iter()
790 .find(|instance| instance.instance_key() == instance_key)
791 }
792
793 pub fn restart_policy_for(&self, instance_key: &str) -> Option<RestartPolicy> {
795 self.module_instance(instance_key)
796 .map(ModuleInstancePlan::restart_policy)
797 }
798
799 pub fn criticality_for(&self, instance_key: &str) -> Option<ModuleCriticality> {
801 self.module_instance(instance_key)
802 .map(ModuleInstancePlan::criticality)
803 }
804
805 pub fn module_instance_is_required(&self, instance_key: &str) -> bool {
807 self.capability_bindings.iter().any(|binding| {
808 binding.provider_instance() == instance_key
809 && self
810 .module_instance(binding.consumer_instance())
811 .is_some_and(|consumer| {
812 consumer.required_capabilities().iter().any(|requirement| {
813 requirement.capability_id() == binding.capability_id()
814 && requirement.cardinality() == CapabilityCardinality::One
815 })
816 })
817 })
818 }
819}