1pub mod authoring;
4
5use std::{
6 collections::{BTreeMap, BTreeSet},
7 fmt,
8 time::Duration,
9};
10
11use serde::{Deserialize, Serialize};
12
13mod execution;
14mod resolution;
15
16pub use execution::{ExecutionClassId, ExecutionLaneId, ExecutionLanePlan};
17use resolution::{
18 activation_order_for, resolve_parts, sort_bindings, sort_module_instances,
19 sorted_execution_lanes, validate_execution_lanes,
20};
21
22pub const PLAN_SCHEMA_VERSION: u32 = 1;
24
25pub const DEFAULT_REQUEST_QUEUE_CAPACITY: usize = 16;
27
28pub const DEFAULT_REQUEST_MAX_CONCURRENCY: usize = 1;
30
31pub const DEFAULT_EVENT_QUEUE_CAPACITY: usize = 16;
33
34fn default_execution_lanes() -> Vec<ExecutionLanePlan> {
35 vec![ExecutionLanePlan::new("main")]
36}
37
38#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
40#[serde(rename_all = "snake_case")]
41pub enum CapabilityCardinality {
42 One,
44 Optional,
46 Many,
48}
49
50#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
52#[serde(rename_all = "snake_case")]
53pub enum CapabilityOperationKind {
54 Request,
56 Stream,
58 Event,
60}
61
62#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
68pub struct RequestAdmissionPlan {
69 queue_capacity: usize,
70 max_concurrency: usize,
71}
72
73impl RequestAdmissionPlan {
74 pub const fn new(queue_capacity: usize, max_concurrency: usize) -> Self {
76 Self {
77 queue_capacity,
78 max_concurrency,
79 }
80 }
81
82 pub const fn queue_capacity(self) -> usize {
84 self.queue_capacity
85 }
86
87 pub const fn max_concurrency(self) -> usize {
89 self.max_concurrency
90 }
91
92 fn validate(self, capability_id: &str, operation: &str) -> Result<(), PlanResolutionError> {
93 if self.max_concurrency == 0 {
94 return Err(PlanResolutionError::InvalidRequestAdmission {
95 capability_id: capability_id.to_owned(),
96 operation: operation.to_owned(),
97 queue_capacity: self.queue_capacity,
98 max_concurrency: self.max_concurrency,
99 });
100 }
101 Ok(())
102 }
103}
104
105impl Default for RequestAdmissionPlan {
106 fn default() -> Self {
107 Self::new(
108 DEFAULT_REQUEST_QUEUE_CAPACITY,
109 DEFAULT_REQUEST_MAX_CONCURRENCY,
110 )
111 }
112}
113
114#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
119pub struct EventAdmissionPlan {
120 capacity: usize,
121}
122
123impl EventAdmissionPlan {
124 pub const fn new(capacity: usize) -> Self {
126 Self { capacity }
127 }
128
129 pub const fn capacity(self) -> usize {
131 self.capacity
132 }
133}
134
135impl Default for EventAdmissionPlan {
136 fn default() -> Self {
137 Self::new(DEFAULT_EVENT_QUEUE_CAPACITY)
138 }
139}
140
141#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
143#[serde(rename_all = "snake_case")]
144pub enum RestartMode {
145 Never,
147 OnFailure,
149}
150
151#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
153pub struct RestartPolicy {
154 mode: RestartMode,
155 max_attempts: usize,
156 window: Duration,
157 backoff: Duration,
158 jitter: Duration,
159 stability: Duration,
160}
161
162impl RestartPolicy {
163 pub const fn never() -> Self {
165 Self {
166 mode: RestartMode::Never,
167 max_attempts: 0,
168 window: Duration::ZERO,
169 backoff: Duration::ZERO,
170 jitter: Duration::ZERO,
171 stability: Duration::ZERO,
172 }
173 }
174
175 pub const fn on_failure(
177 max_attempts: usize,
178 window: Duration,
179 backoff: Duration,
180 jitter: Duration,
181 stability: Duration,
182 ) -> Self {
183 Self {
184 mode: RestartMode::OnFailure,
185 max_attempts,
186 window,
187 backoff,
188 jitter,
189 stability,
190 }
191 }
192
193 pub const fn mode(self) -> RestartMode {
195 self.mode
196 }
197
198 pub const fn max_attempts(self) -> usize {
200 self.max_attempts
201 }
202
203 pub const fn window(self) -> Duration {
205 self.window
206 }
207
208 pub const fn backoff(self) -> Duration {
210 self.backoff
211 }
212
213 pub const fn jitter(self) -> Duration {
215 self.jitter
216 }
217
218 pub const fn stability(self) -> Duration {
220 self.stability
221 }
222
223 fn validate(&self, instance_key: &str) -> Result<(), PlanResolutionError> {
224 if self.mode == RestartMode::OnFailure && (self.max_attempts == 0 || self.window.is_zero())
225 {
226 return Err(PlanResolutionError::InvalidRestartPolicy {
227 instance_key: instance_key.to_owned(),
228 max_attempts: self.max_attempts,
229 window: self.window,
230 });
231 }
232 Ok(())
233 }
234}
235
236impl Default for RestartPolicy {
237 fn default() -> Self {
238 Self::never()
239 }
240}
241
242#[derive(Clone, Copy, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
244#[serde(rename_all = "snake_case")]
245pub enum ModuleCriticality {
246 #[default]
248 NonCritical,
249 Critical,
251}
252
253impl ModuleCriticality {
254 pub const fn is_critical(self) -> bool {
256 matches!(self, Self::Critical)
257 }
258}
259
260#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
262pub struct CapabilityRequirementPlan {
263 capability_id: String,
264 descriptor_version: String,
265 cardinality: CapabilityCardinality,
266}
267
268impl CapabilityRequirementPlan {
269 pub fn new(
271 capability_id: impl Into<String>,
272 descriptor_version: impl Into<String>,
273 cardinality: CapabilityCardinality,
274 ) -> Self {
275 Self {
276 capability_id: capability_id.into(),
277 descriptor_version: descriptor_version.into(),
278 cardinality,
279 }
280 }
281
282 pub fn one(capability_id: impl Into<String>, descriptor_version: impl Into<String>) -> Self {
284 Self::new(
285 capability_id,
286 descriptor_version,
287 CapabilityCardinality::One,
288 )
289 }
290
291 pub fn optional(
293 capability_id: impl Into<String>,
294 descriptor_version: impl Into<String>,
295 ) -> Self {
296 Self::new(
297 capability_id,
298 descriptor_version,
299 CapabilityCardinality::Optional,
300 )
301 }
302
303 pub fn many(capability_id: impl Into<String>, descriptor_version: impl Into<String>) -> Self {
305 Self::new(
306 capability_id,
307 descriptor_version,
308 CapabilityCardinality::Many,
309 )
310 }
311
312 pub fn capability_id(&self) -> &str {
314 &self.capability_id
315 }
316
317 pub fn descriptor_version(&self) -> &str {
319 &self.descriptor_version
320 }
321
322 pub const fn cardinality(&self) -> CapabilityCardinality {
324 self.cardinality
325 }
326}
327
328#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
330pub struct CapabilityEndpointPlan {
331 capability_id: String,
332 descriptor_version: String,
333 operations: Vec<String>,
334 operation_kinds: BTreeMap<String, CapabilityOperationKind>,
335 default_admission: Option<RequestAdmissionPlan>,
336 operation_admissions: BTreeMap<String, RequestAdmissionPlan>,
337 event_admission: Option<EventAdmissionPlan>,
338 #[serde(default)]
339 cross_lane_transfer: bool,
340}
341
342impl CapabilityEndpointPlan {
343 pub fn new(
345 capability_id: impl Into<String>,
346 descriptor_version: impl Into<String>,
347 operations: impl IntoIterator<Item = impl Into<String>>,
348 ) -> Self {
349 Self {
350 capability_id: capability_id.into(),
351 descriptor_version: descriptor_version.into(),
352 operations: operations.into_iter().map(Into::into).collect(),
353 operation_kinds: BTreeMap::new(),
354 default_admission: None,
355 operation_admissions: BTreeMap::new(),
356 event_admission: None,
357 cross_lane_transfer: false,
358 }
359 }
360
361 #[must_use]
363 pub fn with_admission(mut self, admission: RequestAdmissionPlan) -> Self {
364 self.default_admission = Some(admission);
365 self
366 }
367
368 #[must_use]
370 pub fn with_limits(self, queue_capacity: usize, max_concurrency: usize) -> Self {
371 self.with_admission(RequestAdmissionPlan::new(queue_capacity, max_concurrency))
372 }
373
374 #[must_use]
376 pub fn with_operation_kind(
377 mut self,
378 operation: impl Into<String>,
379 kind: CapabilityOperationKind,
380 ) -> Self {
381 self.operation_kinds.insert(operation.into(), kind);
382 self
383 }
384
385 #[must_use]
387 pub fn with_stream_operation(self, operation: impl Into<String>) -> Self {
388 self.with_operation_kind(operation, CapabilityOperationKind::Stream)
389 }
390
391 #[must_use]
393 pub fn with_event_operation(self, operation: impl Into<String>) -> Self {
394 self.with_operation_kind(operation, CapabilityOperationKind::Event)
395 }
396
397 #[must_use]
399 pub fn with_event_admission(mut self, admission: EventAdmissionPlan) -> Self {
400 self.event_admission = Some(admission);
401 self
402 }
403
404 #[must_use]
406 pub fn with_event_capacity(self, capacity: usize) -> Self {
407 self.with_event_admission(EventAdmissionPlan::new(capacity))
408 }
409
410 #[must_use]
412 pub const fn with_cross_lane_transfer(mut self) -> Self {
413 self.cross_lane_transfer = true;
414 self
415 }
416
417 #[must_use]
419 pub fn with_operation_admission(
420 mut self,
421 operation: impl Into<String>,
422 admission: RequestAdmissionPlan,
423 ) -> Self {
424 self.operation_admissions
425 .insert(operation.into(), admission);
426 self
427 }
428
429 #[must_use]
431 pub fn with_operation_limits(
432 self,
433 operation: impl Into<String>,
434 queue_capacity: usize,
435 max_concurrency: usize,
436 ) -> Self {
437 self.with_operation_admission(
438 operation,
439 RequestAdmissionPlan::new(queue_capacity, max_concurrency),
440 )
441 }
442
443 pub fn capability_id(&self) -> &str {
445 &self.capability_id
446 }
447
448 pub fn descriptor_version(&self) -> &str {
450 &self.descriptor_version
451 }
452
453 pub fn operations(&self) -> &[String] {
455 &self.operations
456 }
457
458 pub fn operation_kind(&self, operation: &str) -> Option<CapabilityOperationKind> {
460 self.operations
461 .iter()
462 .any(|declared| declared == operation)
463 .then(|| {
464 self.operation_kinds
465 .get(operation)
466 .copied()
467 .unwrap_or(CapabilityOperationKind::Request)
468 })
469 }
470
471 pub fn stream_operations(&self) -> Vec<&str> {
473 self.operations
474 .iter()
475 .filter(|operation| {
476 self.operation_kind(operation) == Some(CapabilityOperationKind::Stream)
477 })
478 .map(String::as_str)
479 .collect()
480 }
481
482 pub fn request_operations(&self) -> Vec<&str> {
484 self.operations
485 .iter()
486 .filter(|operation| {
487 self.operation_kind(operation) == Some(CapabilityOperationKind::Request)
488 })
489 .map(String::as_str)
490 .collect()
491 }
492
493 pub fn event_operations(&self) -> Vec<&str> {
495 self.operations
496 .iter()
497 .filter(|operation| {
498 self.operation_kind(operation) == Some(CapabilityOperationKind::Event)
499 })
500 .map(String::as_str)
501 .collect()
502 }
503
504 pub const fn event_admission(&self) -> Option<EventAdmissionPlan> {
506 self.event_admission
507 }
508
509 pub const fn supports_cross_lane_transfer(&self) -> bool {
511 self.cross_lane_transfer
512 }
513
514 pub fn default_admission(&self) -> Option<RequestAdmissionPlan> {
516 self.default_admission
517 }
518
519 pub fn operation_admissions(&self) -> &BTreeMap<String, RequestAdmissionPlan> {
521 &self.operation_admissions
522 }
523
524 pub fn operation_admission(&self, operation: &str) -> Option<RequestAdmissionPlan> {
526 self.operation_admissions
527 .get(operation)
528 .copied()
529 .or(self.default_admission)
530 }
531}
532
533#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
535pub struct ModuleInstancePlan {
536 instance_key: String,
537 package_id: String,
538 entrypoint: String,
539 configuration: String,
540 provided_capabilities: Vec<CapabilityEndpointPlan>,
541 required_capabilities: Vec<CapabilityRequirementPlan>,
542 execution_class: ExecutionClassId,
543 package_revision: String,
544 restart_policy: RestartPolicy,
545 criticality: ModuleCriticality,
546 #[serde(default)]
547 execution_lane: ExecutionLaneId,
548}
549
550impl ModuleInstancePlan {
551 pub fn new(instance_key: impl Into<String>, package_id: impl Into<String>) -> Self {
553 Self {
554 instance_key: instance_key.into(),
555 package_id: package_id.into(),
556 entrypoint: "default".to_owned(),
557 configuration: "{}".to_owned(),
558 provided_capabilities: Vec::new(),
559 required_capabilities: Vec::new(),
560 execution_class: ExecutionClassId::native_rust(),
561 package_revision: String::new(),
562 restart_policy: RestartPolicy::default(),
563 criticality: ModuleCriticality::default(),
564 execution_lane: ExecutionLaneId::default(),
565 }
566 }
567
568 #[must_use]
570 pub fn with_entrypoint(mut self, entrypoint: impl Into<String>) -> Self {
571 self.entrypoint = entrypoint.into();
572 self
573 }
574
575 #[must_use]
577 pub fn with_configuration(mut self, configuration: impl Into<String>) -> Self {
578 self.configuration = configuration.into();
579 self
580 }
581
582 #[must_use]
584 pub fn with_capability(mut self, capability: CapabilityEndpointPlan) -> Self {
585 self.provided_capabilities.push(capability);
586 self
587 }
588
589 #[must_use]
591 pub fn with_requirement(mut self, requirement: CapabilityRequirementPlan) -> Self {
592 self.required_capabilities.push(requirement);
593 self
594 }
595
596 #[must_use]
598 pub fn with_required_capability(self, requirement: CapabilityRequirementPlan) -> Self {
599 self.with_requirement(requirement)
600 }
601
602 #[must_use]
604 pub fn with_execution_class(mut self, execution_class: ExecutionClassId) -> Self {
605 self.execution_class = execution_class;
606 self
607 }
608
609 #[must_use]
611 pub fn with_execution_lane(mut self, execution_lane: ExecutionLaneId) -> Self {
612 self.execution_lane = execution_lane;
613 self
614 }
615
616 #[must_use]
618 pub fn with_package_revision(mut self, revision: impl Into<String>) -> Self {
619 self.package_revision = revision.into();
620 self
621 }
622
623 #[must_use]
625 pub fn with_restart_policy(mut self, restart_policy: RestartPolicy) -> Self {
626 self.restart_policy = restart_policy;
627 self
628 }
629
630 #[must_use]
632 pub fn with_criticality(mut self, criticality: ModuleCriticality) -> Self {
633 self.criticality = criticality;
634 self
635 }
636
637 pub fn instance_key(&self) -> &str {
639 &self.instance_key
640 }
641
642 pub fn package_id(&self) -> &str {
644 &self.package_id
645 }
646
647 pub fn entrypoint(&self) -> &str {
649 &self.entrypoint
650 }
651
652 pub fn configuration(&self) -> &str {
654 &self.configuration
655 }
656
657 pub fn provided_capabilities(&self) -> &[CapabilityEndpointPlan] {
659 &self.provided_capabilities
660 }
661
662 pub fn required_capabilities(&self) -> &[CapabilityRequirementPlan] {
664 &self.required_capabilities
665 }
666
667 pub fn execution_class(&self) -> &ExecutionClassId {
669 &self.execution_class
670 }
671
672 pub const fn execution_lane(&self) -> &ExecutionLaneId {
674 &self.execution_lane
675 }
676
677 pub fn package_revision(&self) -> &str {
679 &self.package_revision
680 }
681
682 pub const fn restart_policy(&self) -> RestartPolicy {
684 self.restart_policy
685 }
686
687 pub const fn criticality(&self) -> ModuleCriticality {
689 self.criticality
690 }
691}
692
693#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
695pub struct CapabilityBinding {
696 consumer_instance: String,
697 capability_id: String,
698 descriptor_version: String,
699 provider_instance: String,
700 provider_order: usize,
701 admission: RequestAdmissionPlan,
702 admission_explicit: bool,
703 event_admission: EventAdmissionPlan,
704 event_admission_explicit: bool,
705}
706
707impl CapabilityBinding {
708 pub fn new(
710 consumer_instance: impl Into<String>,
711 capability_id: impl Into<String>,
712 descriptor_version: impl Into<String>,
713 provider_instance: impl Into<String>,
714 ) -> Self {
715 Self {
716 consumer_instance: consumer_instance.into(),
717 capability_id: capability_id.into(),
718 descriptor_version: descriptor_version.into(),
719 provider_instance: provider_instance.into(),
720 provider_order: 0,
721 admission: RequestAdmissionPlan::default(),
722 admission_explicit: false,
723 event_admission: EventAdmissionPlan::default(),
724 event_admission_explicit: false,
725 }
726 }
727
728 #[must_use]
730 pub fn with_admission(mut self, admission: RequestAdmissionPlan) -> Self {
731 self.admission = admission;
732 self.admission_explicit = true;
733 self
734 }
735
736 #[must_use]
738 pub fn with_limits(self, queue_capacity: usize, max_concurrency: usize) -> Self {
739 self.with_admission(RequestAdmissionPlan::new(queue_capacity, max_concurrency))
740 }
741
742 #[must_use]
744 pub fn with_event_admission(mut self, admission: EventAdmissionPlan) -> Self {
745 self.event_admission = admission;
746 self.event_admission_explicit = true;
747 self
748 }
749
750 #[must_use]
752 pub fn with_event_capacity(self, capacity: usize) -> Self {
753 self.with_event_admission(EventAdmissionPlan::new(capacity))
754 }
755
756 fn with_provider_order(mut self, provider_order: usize) -> Self {
757 self.provider_order = provider_order;
758 self
759 }
760
761 pub fn consumer_instance(&self) -> &str {
763 &self.consumer_instance
764 }
765
766 pub fn capability_id(&self) -> &str {
768 &self.capability_id
769 }
770
771 pub fn descriptor_version(&self) -> &str {
773 &self.descriptor_version
774 }
775
776 pub fn provider_instance(&self) -> &str {
778 &self.provider_instance
779 }
780
781 pub const fn provider_order(&self) -> usize {
783 self.provider_order
784 }
785
786 pub const fn admission(&self) -> RequestAdmissionPlan {
788 self.admission
789 }
790
791 pub const fn has_explicit_admission(&self) -> bool {
793 self.admission_explicit
794 }
795
796 pub const fn event_admission(&self) -> EventAdmissionPlan {
798 self.event_admission
799 }
800
801 pub const fn has_explicit_event_admission(&self) -> bool {
803 self.event_admission_explicit
804 }
805}
806
807#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
809pub struct AppComposition {
810 module_instances: Vec<ModuleInstancePlan>,
811 capability_bindings: Vec<CapabilityBinding>,
812 #[serde(default = "default_execution_lanes")]
813 execution_lanes: Vec<ExecutionLanePlan>,
814}
815
816impl AppComposition {
817 pub fn new(
819 module_instances: Vec<ModuleInstancePlan>,
820 capability_bindings: Vec<CapabilityBinding>,
821 ) -> Self {
822 Self {
823 module_instances,
824 capability_bindings,
825 execution_lanes: default_execution_lanes(),
826 }
827 }
828
829 #[must_use]
831 pub fn with_execution_lanes(mut self, execution_lanes: Vec<ExecutionLanePlan>) -> Self {
832 self.execution_lanes = execution_lanes;
833 self
834 }
835
836 pub fn resolve(&self) -> Result<ResolvedAppPlan, PlanResolutionError> {
838 validate_execution_lanes(&self.execution_lanes, &self.module_instances)?;
839 resolve_parts(&self.module_instances, &self.capability_bindings).map(
840 |(module_instances, capability_bindings)| ResolvedAppPlan {
841 schema_version: PLAN_SCHEMA_VERSION,
842 module_instances,
843 capability_bindings,
844 execution_lanes: sorted_execution_lanes(&self.execution_lanes),
845 },
846 )
847 }
848
849 pub fn module_instances(&self) -> &[ModuleInstancePlan] {
851 &self.module_instances
852 }
853
854 pub fn capability_bindings(&self) -> &[CapabilityBinding] {
856 &self.capability_bindings
857 }
858
859 pub fn execution_lanes(&self) -> &[ExecutionLanePlan] {
861 &self.execution_lanes
862 }
863}
864
865#[derive(Clone, Debug, Eq, PartialEq)]
867pub enum PlanResolutionError {
868 UnsupportedSchemaVersion { expected: u32, actual: u32 },
870 DuplicateModuleInstance { instance_key: String },
872 MissingExecutionLane,
874 InvalidExecutionLane { execution_lane: String },
876 DuplicateExecutionLane { execution_lane: String },
878 UndeclaredExecutionLane {
880 instance_key: String,
881 execution_lane: String,
882 },
883 InvalidModuleEntrypoint { instance_key: String },
885 DuplicateProvidedCapability {
887 provider_instance: String,
888 capability_id: String,
889 },
890 DuplicateOperation {
892 provider_instance: String,
893 capability_id: String,
894 operation: String,
895 },
896 DuplicateRequiredCapability {
898 consumer_instance: String,
899 capability_id: String,
900 },
901 InvalidConsumerReference {
903 consumer_instance: String,
904 capability_id: String,
905 },
906 InvalidProviderReference {
908 consumer_instance: String,
909 capability_id: String,
910 provider_instance: String,
911 },
912 UndeclaredCapabilityRequirement {
914 consumer_instance: String,
915 capability_id: String,
916 },
917 IncompatibleCapabilityVersion {
919 consumer_instance: String,
920 capability_id: String,
921 required: String,
922 provided: String,
923 provider_instance: String,
924 },
925 CrossLaneTransferUnsupported {
927 consumer_instance: String,
928 provider_instance: String,
929 capability_id: String,
930 },
931 CrossLaneInteractionUnsupported {
933 capability_id: String,
934 operation: String,
935 interaction: CapabilityOperationKind,
936 },
937 MissingOneBinding {
939 consumer_instance: String,
940 capability_id: String,
941 },
942 AmbiguousOneBinding {
944 consumer_instance: String,
945 capability_id: String,
946 providers: usize,
947 },
948 AmbiguousOptionalBinding {
950 consumer_instance: String,
951 capability_id: String,
952 providers: usize,
953 },
954 DuplicateBinding {
956 consumer_instance: String,
957 capability_id: String,
958 provider_instance: String,
959 },
960 InvalidRequestAdmission {
962 capability_id: String,
963 operation: String,
964 queue_capacity: usize,
965 max_concurrency: usize,
966 },
967 UnknownAdmissionOperation {
969 capability_id: String,
970 operation: String,
971 },
972 UnknownOperationInteraction {
974 capability_id: String,
975 operation: String,
976 },
977 InvalidRestartPolicy {
979 instance_key: String,
980 max_attempts: usize,
981 window: Duration,
982 },
983 ActivationCycle { instances: Vec<String> },
985}
986
987impl fmt::Display for PlanResolutionError {
988 #[allow(clippy::too_many_lines)]
989 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
990 match self {
991 Self::UnsupportedSchemaVersion { expected, actual } => write!(
992 formatter,
993 "unsupported Plan schema version {actual}; expected {expected}"
994 ),
995 Self::DuplicateModuleInstance { instance_key } => {
996 write!(formatter, "duplicate Module Instance `{instance_key}`")
997 }
998 Self::MissingExecutionLane => {
999 formatter.write_str("Resolved App Plan declares no Execution Lanes")
1000 }
1001 Self::InvalidExecutionLane { execution_lane } => {
1002 write!(formatter, "invalid Execution Lane `{execution_lane}`")
1003 }
1004 Self::DuplicateExecutionLane { execution_lane } => {
1005 write!(formatter, "duplicate Execution Lane `{execution_lane}`")
1006 }
1007 Self::UndeclaredExecutionLane {
1008 instance_key,
1009 execution_lane,
1010 } => write!(
1011 formatter,
1012 "Module Instance `{instance_key}` is placed on undeclared Execution Lane `{execution_lane}`"
1013 ),
1014 Self::InvalidModuleEntrypoint { instance_key } => write!(
1015 formatter,
1016 "Module Instance `{instance_key}` has an empty entrypoint"
1017 ),
1018 Self::DuplicateProvidedCapability {
1019 provider_instance,
1020 capability_id,
1021 } => write!(
1022 formatter,
1023 "Module Instance `{provider_instance}` provides Capability `{capability_id}` more than once"
1024 ),
1025 Self::DuplicateOperation {
1026 provider_instance,
1027 capability_id,
1028 operation,
1029 } => write!(
1030 formatter,
1031 "Module Instance `{provider_instance}` Capability `{capability_id}` declares Operation `{operation}` more than once"
1032 ),
1033 Self::DuplicateRequiredCapability {
1034 consumer_instance,
1035 capability_id,
1036 } => write!(
1037 formatter,
1038 "Module Instance `{consumer_instance}` requires Capability `{capability_id}` more than once"
1039 ),
1040 Self::InvalidConsumerReference {
1041 consumer_instance,
1042 capability_id,
1043 } => write!(
1044 formatter,
1045 "Capability `{capability_id}` names missing consumer `{consumer_instance}`"
1046 ),
1047 Self::InvalidProviderReference {
1048 consumer_instance,
1049 capability_id,
1050 provider_instance,
1051 } => write!(
1052 formatter,
1053 "consumer `{consumer_instance}` names invalid provider `{provider_instance}` for Capability `{capability_id}`"
1054 ),
1055 Self::UndeclaredCapabilityRequirement {
1056 consumer_instance,
1057 capability_id,
1058 } => write!(
1059 formatter,
1060 "consumer `{consumer_instance}` has no declared requirement for Capability `{capability_id}`"
1061 ),
1062 Self::IncompatibleCapabilityVersion {
1063 consumer_instance,
1064 capability_id,
1065 required,
1066 provided,
1067 provider_instance,
1068 } => write!(
1069 formatter,
1070 "consumer `{consumer_instance}` requires Capability `{capability_id}` version `{required}`, but provider `{provider_instance}` provides `{provided}`"
1071 ),
1072 Self::CrossLaneTransferUnsupported {
1073 consumer_instance,
1074 provider_instance,
1075 capability_id,
1076 } => write!(
1077 formatter,
1078 "consumer `{consumer_instance}` binds Capability `{capability_id}` across Execution Lanes to provider `{provider_instance}`, but its contract types do not support cross-lane transfer"
1079 ),
1080 Self::CrossLaneInteractionUnsupported {
1081 capability_id,
1082 operation,
1083 interaction,
1084 } => write!(
1085 formatter,
1086 "Capability `{capability_id}` Operation `{operation}` uses {interaction:?}, which this Plan version cannot transfer across Execution Lanes"
1087 ),
1088 Self::MissingOneBinding {
1089 consumer_instance,
1090 capability_id,
1091 } => write!(
1092 formatter,
1093 "consumer `{consumer_instance}` is missing one binding for Capability `{capability_id}`"
1094 ),
1095 Self::AmbiguousOneBinding {
1096 consumer_instance,
1097 capability_id,
1098 providers,
1099 } => write!(
1100 formatter,
1101 "consumer `{consumer_instance}` has {providers} bindings for one Capability `{capability_id}`"
1102 ),
1103 Self::AmbiguousOptionalBinding {
1104 consumer_instance,
1105 capability_id,
1106 providers,
1107 } => write!(
1108 formatter,
1109 "consumer `{consumer_instance}` has {providers} bindings for optional Capability `{capability_id}`"
1110 ),
1111 Self::DuplicateBinding {
1112 consumer_instance,
1113 capability_id,
1114 provider_instance,
1115 } => write!(
1116 formatter,
1117 "consumer `{consumer_instance}` binds Capability `{capability_id}` to provider `{provider_instance}` more than once"
1118 ),
1119 Self::InvalidRequestAdmission {
1120 capability_id,
1121 operation,
1122 queue_capacity,
1123 max_concurrency,
1124 } => write!(
1125 formatter,
1126 "Capability `{capability_id}` Operation `{operation}` has invalid request admission (queue capacity {queue_capacity}, concurrency {max_concurrency})"
1127 ),
1128 Self::UnknownAdmissionOperation {
1129 capability_id,
1130 operation,
1131 } => write!(
1132 formatter,
1133 "Capability `{capability_id}` configures request admission for unknown Operation `{operation}`"
1134 ),
1135 Self::UnknownOperationInteraction {
1136 capability_id,
1137 operation,
1138 } => write!(
1139 formatter,
1140 "Capability `{capability_id}` configures interaction metadata for unknown Operation `{operation}`"
1141 ),
1142 Self::InvalidRestartPolicy {
1143 instance_key,
1144 max_attempts,
1145 window,
1146 } => write!(
1147 formatter,
1148 "Module Instance `{instance_key}` has invalid restart policy (attempts {max_attempts}, window {window:?})"
1149 ),
1150 Self::ActivationCycle { instances } => write!(
1151 formatter,
1152 "required Capability activation cycle: {}",
1153 instances.join(" -> ")
1154 ),
1155 }
1156 }
1157}
1158
1159impl std::error::Error for PlanResolutionError {}
1160
1161#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
1163pub struct ResolvedAppPlan {
1164 schema_version: u32,
1165 module_instances: Vec<ModuleInstancePlan>,
1166 capability_bindings: Vec<CapabilityBinding>,
1167 #[serde(default = "default_execution_lanes")]
1168 execution_lanes: Vec<ExecutionLanePlan>,
1169}
1170
1171impl ResolvedAppPlan {
1172 pub fn empty() -> Self {
1174 Self {
1175 schema_version: PLAN_SCHEMA_VERSION,
1176 module_instances: Vec::new(),
1177 capability_bindings: Vec::new(),
1178 execution_lanes: default_execution_lanes(),
1179 }
1180 }
1181
1182 pub fn new(
1184 mut module_instances: Vec<ModuleInstancePlan>,
1185 mut capability_bindings: Vec<CapabilityBinding>,
1186 ) -> Self {
1187 sort_module_instances(&mut module_instances);
1188 sort_bindings(&mut capability_bindings);
1189 Self {
1190 schema_version: PLAN_SCHEMA_VERSION,
1191 module_instances,
1192 capability_bindings,
1193 execution_lanes: default_execution_lanes(),
1194 }
1195 }
1196
1197 pub const fn with_schema_version(schema_version: u32) -> Self {
1201 Self {
1202 schema_version,
1203 module_instances: Vec::new(),
1204 capability_bindings: Vec::new(),
1205 execution_lanes: Vec::new(),
1206 }
1207 }
1208
1209 pub fn validate(&self) -> Result<(), PlanResolutionError> {
1211 if self.schema_version != PLAN_SCHEMA_VERSION {
1212 return Err(PlanResolutionError::UnsupportedSchemaVersion {
1213 expected: PLAN_SCHEMA_VERSION,
1214 actual: self.schema_version,
1215 });
1216 }
1217 validate_execution_lanes(&self.execution_lanes, &self.module_instances)?;
1218 resolve_parts(&self.module_instances, &self.capability_bindings).map(|_| ())
1219 }
1220
1221 pub fn activation_order(&self) -> Result<Vec<String>, PlanResolutionError> {
1226 if self.schema_version != PLAN_SCHEMA_VERSION {
1227 return Err(PlanResolutionError::UnsupportedSchemaVersion {
1228 expected: PLAN_SCHEMA_VERSION,
1229 actual: self.schema_version,
1230 });
1231 }
1232 validate_execution_lanes(&self.execution_lanes, &self.module_instances)?;
1233 let (instances, bindings) =
1234 resolve_parts(&self.module_instances, &self.capability_bindings)?;
1235 activation_order_for(&instances, &bindings)
1236 .map_err(|instances| PlanResolutionError::ActivationCycle { instances })
1237 }
1238
1239 pub const fn schema_version(&self) -> u32 {
1241 self.schema_version
1242 }
1243
1244 pub fn module_instances(&self) -> &[ModuleInstancePlan] {
1246 &self.module_instances
1247 }
1248
1249 pub fn capability_bindings(&self) -> &[CapabilityBinding] {
1251 &self.capability_bindings
1252 }
1253
1254 pub fn execution_lanes(&self) -> &[ExecutionLanePlan] {
1256 &self.execution_lanes
1257 }
1258
1259 pub fn request_admission_for(
1261 &self,
1262 binding: &CapabilityBinding,
1263 operation: &str,
1264 ) -> RequestAdmissionPlan {
1265 if binding.has_explicit_admission() {
1266 return binding.admission();
1267 }
1268
1269 self.module_instances
1270 .iter()
1271 .find(|instance| instance.instance_key() == binding.provider_instance())
1272 .and_then(|provider| {
1273 provider
1274 .provided_capabilities()
1275 .iter()
1276 .find(|endpoint| endpoint.capability_id() == binding.capability_id())
1277 })
1278 .and_then(|endpoint| endpoint.operation_admission(operation))
1279 .unwrap_or_else(|| binding.admission())
1280 }
1281
1282 pub fn event_admission_for(&self, binding: &CapabilityBinding) -> EventAdmissionPlan {
1284 if binding.has_explicit_event_admission() {
1285 return binding.event_admission();
1286 }
1287
1288 self.module_instances
1289 .iter()
1290 .find(|instance| instance.instance_key() == binding.provider_instance())
1291 .and_then(|provider| {
1292 provider
1293 .provided_capabilities()
1294 .iter()
1295 .find(|endpoint| endpoint.capability_id() == binding.capability_id())
1296 })
1297 .and_then(CapabilityEndpointPlan::event_admission)
1298 .unwrap_or_else(|| binding.event_admission())
1299 }
1300
1301 pub fn module_instance(&self, instance_key: &str) -> Option<&ModuleInstancePlan> {
1303 self.module_instances
1304 .iter()
1305 .find(|instance| instance.instance_key() == instance_key)
1306 }
1307
1308 pub fn restart_policy_for(&self, instance_key: &str) -> Option<RestartPolicy> {
1310 self.module_instance(instance_key)
1311 .map(ModuleInstancePlan::restart_policy)
1312 }
1313
1314 pub fn criticality_for(&self, instance_key: &str) -> Option<ModuleCriticality> {
1316 self.module_instance(instance_key)
1317 .map(ModuleInstancePlan::criticality)
1318 }
1319
1320 pub fn module_instance_is_required(&self, instance_key: &str) -> bool {
1322 self.capability_bindings.iter().any(|binding| {
1323 binding.provider_instance() == instance_key
1324 && self
1325 .module_instance(binding.consumer_instance())
1326 .is_some_and(|consumer| {
1327 consumer.required_capabilities().iter().any(|requirement| {
1328 requirement.capability_id() == binding.capability_id()
1329 && requirement.cardinality() == CapabilityCardinality::One
1330 })
1331 })
1332 })
1333 }
1334}