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 PLUGIN_AUTHORING_V2_RUNTIME_PROFILE: &str = "lenso.plugin-authoring@2";
34
35pub const DEFAULT_REQUEST_QUEUE_CAPACITY: usize = 16;
37
38pub const DEFAULT_REQUEST_MAX_CONCURRENCY: usize = 1;
40
41pub const DEFAULT_EVENT_QUEUE_CAPACITY: usize = 16;
43
44fn default_execution_lanes() -> Vec<ExecutionLanePlan> {
45 vec![ExecutionLanePlan::new("main")]
46}
47
48#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
50pub struct CapabilityEndpointPlan {
51 capability_id: String,
52 descriptor_version: String,
53 operations: Vec<String>,
54 operation_kinds: BTreeMap<String, CapabilityOperationKind>,
55 default_admission: Option<RequestAdmissionPlan>,
56 operation_admissions: BTreeMap<String, RequestAdmissionPlan>,
57 event_admission: Option<EventAdmissionPlan>,
58 #[serde(default)]
59 cross_lane_transfer: bool,
60}
61
62impl CapabilityEndpointPlan {
63 pub fn new(
65 capability_id: impl Into<String>,
66 descriptor_version: impl Into<String>,
67 operations: impl IntoIterator<Item = impl Into<String>>,
68 ) -> Self {
69 Self {
70 capability_id: capability_id.into(),
71 descriptor_version: descriptor_version.into(),
72 operations: operations.into_iter().map(Into::into).collect(),
73 operation_kinds: BTreeMap::new(),
74 default_admission: None,
75 operation_admissions: BTreeMap::new(),
76 event_admission: None,
77 cross_lane_transfer: false,
78 }
79 }
80
81 #[must_use]
83 pub fn with_admission(mut self, admission: RequestAdmissionPlan) -> Self {
84 self.default_admission = Some(admission);
85 self
86 }
87
88 #[must_use]
90 pub fn with_limits(self, queue_capacity: usize, max_concurrency: usize) -> Self {
91 self.with_admission(RequestAdmissionPlan::new(queue_capacity, max_concurrency))
92 }
93
94 #[must_use]
96 pub fn with_operation_kind(
97 mut self,
98 operation: impl Into<String>,
99 kind: CapabilityOperationKind,
100 ) -> Self {
101 self.operation_kinds.insert(operation.into(), kind);
102 self
103 }
104
105 #[must_use]
107 pub fn with_stream_operation(self, operation: impl Into<String>) -> Self {
108 self.with_operation_kind(operation, CapabilityOperationKind::Stream)
109 }
110
111 #[must_use]
113 pub fn with_event_operation(self, operation: impl Into<String>) -> Self {
114 self.with_operation_kind(operation, CapabilityOperationKind::Event)
115 }
116
117 #[must_use]
119 pub fn with_event_admission(mut self, admission: EventAdmissionPlan) -> Self {
120 self.event_admission = Some(admission);
121 self
122 }
123
124 #[must_use]
126 pub fn with_event_capacity(self, capacity: usize) -> Self {
127 self.with_event_admission(EventAdmissionPlan::new(capacity))
128 }
129
130 #[must_use]
132 pub const fn with_cross_lane_transfer(mut self) -> Self {
133 self.cross_lane_transfer = true;
134 self
135 }
136
137 #[must_use]
139 pub fn with_operation_admission(
140 mut self,
141 operation: impl Into<String>,
142 admission: RequestAdmissionPlan,
143 ) -> Self {
144 self.operation_admissions
145 .insert(operation.into(), admission);
146 self
147 }
148
149 #[must_use]
151 pub fn with_operation_limits(
152 self,
153 operation: impl Into<String>,
154 queue_capacity: usize,
155 max_concurrency: usize,
156 ) -> Self {
157 self.with_operation_admission(
158 operation,
159 RequestAdmissionPlan::new(queue_capacity, max_concurrency),
160 )
161 }
162
163 pub fn capability_id(&self) -> &str {
165 &self.capability_id
166 }
167
168 pub fn descriptor_version(&self) -> &str {
170 &self.descriptor_version
171 }
172
173 pub fn operations(&self) -> &[String] {
175 &self.operations
176 }
177
178 pub fn operation_kind(&self, operation: &str) -> Option<CapabilityOperationKind> {
180 self.operations
181 .iter()
182 .any(|declared| declared == operation)
183 .then(|| {
184 self.operation_kinds
185 .get(operation)
186 .copied()
187 .unwrap_or(CapabilityOperationKind::Request)
188 })
189 }
190
191 pub fn stream_operations(&self) -> Vec<&str> {
193 self.operations
194 .iter()
195 .filter(|operation| {
196 self.operation_kind(operation) == Some(CapabilityOperationKind::Stream)
197 })
198 .map(String::as_str)
199 .collect()
200 }
201
202 pub fn request_operations(&self) -> Vec<&str> {
204 self.operations
205 .iter()
206 .filter(|operation| {
207 self.operation_kind(operation) == Some(CapabilityOperationKind::Request)
208 })
209 .map(String::as_str)
210 .collect()
211 }
212
213 pub fn event_operations(&self) -> Vec<&str> {
215 self.operations
216 .iter()
217 .filter(|operation| {
218 self.operation_kind(operation) == Some(CapabilityOperationKind::Event)
219 })
220 .map(String::as_str)
221 .collect()
222 }
223
224 pub const fn event_admission(&self) -> Option<EventAdmissionPlan> {
226 self.event_admission
227 }
228
229 pub const fn supports_cross_lane_transfer(&self) -> bool {
231 self.cross_lane_transfer
232 }
233
234 pub fn default_admission(&self) -> Option<RequestAdmissionPlan> {
236 self.default_admission
237 }
238
239 pub fn operation_admissions(&self) -> &BTreeMap<String, RequestAdmissionPlan> {
241 &self.operation_admissions
242 }
243
244 pub fn operation_admission(&self, operation: &str) -> Option<RequestAdmissionPlan> {
246 self.operation_admissions
247 .get(operation)
248 .copied()
249 .or(self.default_admission)
250 }
251}
252
253#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
255pub struct AppComposition {
256 plugin_instances: Vec<PluginInstancePlan>,
257 capability_bindings: Vec<CapabilityBinding>,
258 #[serde(default = "default_execution_lanes")]
259 execution_lanes: Vec<ExecutionLanePlan>,
260}
261
262impl AppComposition {
263 pub fn new(
265 plugin_instances: Vec<PluginInstancePlan>,
266 capability_bindings: Vec<CapabilityBinding>,
267 ) -> Self {
268 Self {
269 plugin_instances,
270 capability_bindings,
271 execution_lanes: default_execution_lanes(),
272 }
273 }
274
275 #[must_use]
277 pub fn with_execution_lanes(mut self, execution_lanes: Vec<ExecutionLanePlan>) -> Self {
278 self.execution_lanes = execution_lanes;
279 self
280 }
281
282 pub fn resolve(&self) -> Result<ResolvedAppPlan, PlanResolutionError> {
284 validate_execution_lanes(&self.execution_lanes, &self.plugin_instances)?;
285 resolve_parts(&self.plugin_instances, &self.capability_bindings).map(
286 |(plugin_instances, capability_bindings)| ResolvedAppPlan {
287 terminal_policy: TerminalPolicy::RequiredPath,
288 schema_version: PLAN_SCHEMA_VERSION,
289 plugin_instances,
290 capability_bindings,
291 execution_lanes: sorted_execution_lanes(&self.execution_lanes),
292 },
293 )
294 }
295
296 pub fn plugin_instances(&self) -> &[PluginInstancePlan] {
298 &self.plugin_instances
299 }
300
301 pub fn capability_bindings(&self) -> &[CapabilityBinding] {
303 &self.capability_bindings
304 }
305
306 pub fn execution_lanes(&self) -> &[ExecutionLanePlan] {
308 &self.execution_lanes
309 }
310}
311
312#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
314#[serde(try_from = "schema::PlanWire")]
315pub struct ResolvedAppPlan {
316 terminal_policy: TerminalPolicy,
317 schema_version: u32,
318 plugin_instances: Vec<PluginInstancePlan>,
319 capability_bindings: Vec<CapabilityBinding>,
320 #[serde(default = "default_execution_lanes")]
321 execution_lanes: Vec<ExecutionLanePlan>,
322}
323
324impl ResolvedAppPlan {
325 pub fn empty() -> Self {
327 Self {
328 terminal_policy: TerminalPolicy::RequiredPath,
329 schema_version: PLAN_SCHEMA_VERSION,
330 plugin_instances: Vec::new(),
331 capability_bindings: Vec::new(),
332 execution_lanes: default_execution_lanes(),
333 }
334 }
335
336 pub fn new(
338 mut plugin_instances: Vec<PluginInstancePlan>,
339 mut capability_bindings: Vec<CapabilityBinding>,
340 ) -> Self {
341 sort_plugin_instances(&mut plugin_instances);
342 sort_bindings(&mut capability_bindings);
343 Self {
344 terminal_policy: TerminalPolicy::RequiredPath,
345 schema_version: PLAN_SCHEMA_VERSION,
346 plugin_instances,
347 capability_bindings,
348 execution_lanes: default_execution_lanes(),
349 }
350 }
351
352 pub const fn with_schema_version(schema_version: u32) -> Self {
356 Self {
357 terminal_policy: TerminalPolicy::RequiredPath,
358 schema_version,
359 plugin_instances: Vec::new(),
360 capability_bindings: Vec::new(),
361 execution_lanes: Vec::new(),
362 }
363 }
364
365 pub fn validate(&self) -> Result<(), PlanResolutionError> {
367 if self.schema_version != PLAN_SCHEMA_VERSION {
368 return Err(PlanResolutionError::UnsupportedSchemaVersion {
369 expected: PLAN_SCHEMA_VERSION,
370 actual: self.schema_version,
371 });
372 }
373 validate_execution_lanes(&self.execution_lanes, &self.plugin_instances)?;
374 let (instances, bindings) =
375 resolve_parts(&self.plugin_instances, &self.capability_bindings)?;
376 self.terminal_policy.validate(&instances, &bindings)
377 }
378
379 pub fn activation_order(&self) -> Result<Vec<String>, PlanResolutionError> {
384 if self.schema_version != PLAN_SCHEMA_VERSION {
385 return Err(PlanResolutionError::UnsupportedSchemaVersion {
386 expected: PLAN_SCHEMA_VERSION,
387 actual: self.schema_version,
388 });
389 }
390 validate_execution_lanes(&self.execution_lanes, &self.plugin_instances)?;
391 let (instances, bindings) =
392 resolve_parts(&self.plugin_instances, &self.capability_bindings)?;
393 self.terminal_policy.validate(&instances, &bindings)?;
394 activation_order_for(&instances, &bindings)
395 .map_err(|instances| PlanResolutionError::ActivationCycle { instances })
396 }
397
398 pub const fn terminal_policy(&self) -> &TerminalPolicy {
400 &self.terminal_policy
401 }
402
403 #[must_use]
405 pub fn with_terminal_policy(mut self, policy: TerminalPolicy) -> Self {
406 self.terminal_policy = policy;
407 self
408 }
409
410 pub const fn schema_version(&self) -> u32 {
412 self.schema_version
413 }
414
415 pub fn plugin_instances(&self) -> &[PluginInstancePlan] {
417 &self.plugin_instances
418 }
419
420 pub fn capability_bindings(&self) -> &[CapabilityBinding] {
422 &self.capability_bindings
423 }
424
425 pub fn execution_lanes(&self) -> &[ExecutionLanePlan] {
427 &self.execution_lanes
428 }
429
430 pub fn request_admission_for(
432 &self,
433 binding: &CapabilityBinding,
434 operation: &str,
435 ) -> RequestAdmissionPlan {
436 if binding.has_explicit_admission() {
437 return binding.admission();
438 }
439
440 self.plugin_instances
441 .iter()
442 .find(|instance| instance.instance_key() == binding.provider_instance())
443 .and_then(|provider| {
444 provider
445 .provided_capabilities()
446 .iter()
447 .find(|endpoint| endpoint.capability_id() == binding.capability_id())
448 })
449 .and_then(|endpoint| endpoint.operation_admission(operation))
450 .unwrap_or_else(|| binding.admission())
451 }
452
453 pub fn event_admission_for(&self, binding: &CapabilityBinding) -> EventAdmissionPlan {
455 if binding.has_explicit_event_admission() {
456 return binding.event_admission();
457 }
458
459 self.plugin_instances
460 .iter()
461 .find(|instance| instance.instance_key() == binding.provider_instance())
462 .and_then(|provider| {
463 provider
464 .provided_capabilities()
465 .iter()
466 .find(|endpoint| endpoint.capability_id() == binding.capability_id())
467 })
468 .and_then(CapabilityEndpointPlan::event_admission)
469 .unwrap_or_else(|| binding.event_admission())
470 }
471
472 pub fn plugin_instance(&self, instance_key: &str) -> Option<&PluginInstancePlan> {
474 self.plugin_instances
475 .iter()
476 .find(|instance| instance.instance_key() == instance_key)
477 }
478
479 pub fn restart_policy_for(&self, instance_key: &str) -> Option<RestartPolicy> {
481 self.plugin_instance(instance_key)
482 .map(PluginInstancePlan::restart_policy)
483 }
484
485 pub fn criticality_for(&self, instance_key: &str) -> Option<PluginCriticality> {
487 self.plugin_instance(instance_key)
488 .map(PluginInstancePlan::criticality)
489 }
490
491 pub fn plugin_instance_is_required(&self, instance_key: &str) -> bool {
493 self.capability_bindings.iter().any(|binding| {
494 binding.provider_instance() == instance_key
495 && self
496 .plugin_instance(binding.consumer_instance())
497 .is_some_and(|consumer| {
498 consumer.required_capabilities().iter().any(|requirement| {
499 requirement.requirement_id() == binding.requirement_id()
500 && requirement.cardinality() == CapabilityCardinality::One
501 })
502 })
503 })
504 }
505
506 pub fn plugin_instance_is_terminal(&self, instance_key: &str) -> bool {
509 match &self.terminal_policy {
510 TerminalPolicy::RequiredPath => self.plugin_instance_is_required(instance_key),
511 TerminalPolicy::HostEssential { closure, .. } => closure
512 .binary_search_by(|candidate| candidate.as_str().cmp(instance_key))
513 .is_ok(),
514 }
515 }
516}