Skip to main content

lenso_app_plan/
lib.rs

1//! App Composition and immutable execution input for the Lenso vNext Kernel.
2
3pub 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
29/// The Resolved App Plan schema understood by this Kernel version.
30pub const PLAN_SCHEMA_VERSION: u32 = 3;
31
32/// Default maximum number of requests waiting for one Operation.
33pub const DEFAULT_REQUEST_QUEUE_CAPACITY: usize = 16;
34
35/// Default maximum concurrent executions for one Operation.
36pub const DEFAULT_REQUEST_MAX_CONCURRENCY: usize = 1;
37
38/// Default maximum number of accepted Events retained by one explicit binding.
39pub const DEFAULT_EVENT_QUEUE_CAPACITY: usize = 16;
40
41fn default_execution_lanes() -> Vec<ExecutionLanePlan> {
42    vec![ExecutionLanePlan::new("main")]
43}
44
45/// Exact Capability endpoint metadata expected from one Plugin Instance.
46#[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    /// Declares one exact Capability Descriptor and its stable Operation table.
61    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    /// Applies one bounded admission policy to every Operation on this endpoint.
79    #[must_use]
80    pub fn with_admission(mut self, admission: RequestAdmissionPlan) -> Self {
81        self.default_admission = Some(admission);
82        self
83    }
84
85    /// Applies queue and concurrency limits to every Operation on this endpoint.
86    #[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    /// Marks one declared Operation with its transport-independent interaction kind.
92    #[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    /// Marks one declared Operation as a bidirectional stream.
103    #[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    /// Marks one declared Operation as an ephemeral Event.
109    #[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    /// Applies one volatile mailbox policy to every Event Operation on this endpoint.
115    #[must_use]
116    pub fn with_event_admission(mut self, admission: EventAdmissionPlan) -> Self {
117        self.event_admission = Some(admission);
118        self
119    }
120
121    /// Applies one volatile mailbox capacity to every Event Operation on this endpoint.
122    #[must_use]
123    pub fn with_event_capacity(self, capacity: usize) -> Self {
124        self.with_event_admission(EventAdmissionPlan::new(capacity))
125    }
126
127    /// Marks the generated Capability value types as safe for native cross-lane transfer.
128    #[must_use]
129    pub const fn with_cross_lane_transfer(mut self) -> Self {
130        self.cross_lane_transfer = true;
131        self
132    }
133
134    /// Applies one bounded admission policy to a named Operation.
135    #[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    /// Applies queue and concurrency limits to a named Operation.
147    #[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    /// Returns the Capability series identity.
161    pub fn capability_id(&self) -> &str {
162        &self.capability_id
163    }
164
165    /// Returns the exact Descriptor version.
166    pub fn descriptor_version(&self) -> &str {
167        &self.descriptor_version
168    }
169
170    /// Returns the exact stable Operation table.
171    pub fn operations(&self) -> &[String] {
172        &self.operations
173    }
174
175    /// Returns the interaction kind of one declared Operation.
176    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    /// Returns the declared stream Operation names in Descriptor order.
189    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    /// Returns the declared request Operation names in Descriptor order.
200    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    /// Returns the declared ephemeral Event Operation names in Descriptor order.
211    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    /// Returns the endpoint-wide Event mailbox policy, when one was authored.
222    pub const fn event_admission(&self) -> Option<EventAdmissionPlan> {
223        self.event_admission
224    }
225
226    /// Returns whether generated values may cross native Execution Lanes without serialization.
227    pub const fn supports_cross_lane_transfer(&self) -> bool {
228        self.cross_lane_transfer
229    }
230
231    /// Returns the endpoint-wide admission policy, when one was authored.
232    pub fn default_admission(&self) -> Option<RequestAdmissionPlan> {
233        self.default_admission
234    }
235
236    /// Returns the Operation-specific admission policies.
237    pub fn operation_admissions(&self) -> &BTreeMap<String, RequestAdmissionPlan> {
238        &self.operation_admissions
239    }
240
241    /// Returns the effective policy for one Operation, if one was authored.
242    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/// Declarative, language-independent authoring input for one App.
251#[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    /// Creates an App Composition with explicit Plugin Instances and bindings.
261    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    /// Replaces the declared Execution Lane set.
273    #[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    /// Materializes one deterministic, validated Resolved App Plan.
280    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    /// Returns the authoring Plugin Instances.
294    pub fn plugin_instances(&self) -> &[PluginInstancePlan] {
295        &self.plugin_instances
296    }
297
298    /// Returns the authoring bindings.
299    pub fn capability_bindings(&self) -> &[CapabilityBinding] {
300        &self.capability_bindings
301    }
302
303    /// Returns the authored Execution Lanes.
304    pub fn execution_lanes(&self) -> &[ExecutionLanePlan] {
305        &self.execution_lanes
306    }
307}
308
309/// Exact, immutable execution input supplied to the Kernel.
310#[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    /// Creates a valid Plan containing no Plugin Instances.
323    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    /// Creates a Plan with exact entries, retaining invalid entries for later validation.
334    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    /// Creates a Plan with an explicit schema version.
350    ///
351    /// This is primarily useful to decode authoring-tool output before validation.
352    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    /// Validates the immutable Plan graph before a Runtime Driver or Adapter boots it.
363    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    /// Returns the deterministic provider-before-consumer lifecycle order.
377    ///
378    /// Every explicit binding is an activation dependency, including an
379    /// optional or many binding when one is present in the resolved Plan.
380    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    /// Returns the Plan schema version.
396    pub const fn terminal_policy(&self) -> &TerminalPolicy {
397        &self.terminal_policy
398    }
399
400    /// Selects an explicit terminal policy; unsupported policies fail validation.
401    #[must_use]
402    pub fn with_terminal_policy(mut self, policy: TerminalPolicy) -> Self {
403        self.terminal_policy = policy;
404        self
405    }
406
407    /// Returns the Plan schema version.
408    pub const fn schema_version(&self) -> u32 {
409        self.schema_version
410    }
411
412    /// Returns the exact Plugin Instances in deterministic Plan order.
413    pub fn plugin_instances(&self) -> &[PluginInstancePlan] {
414        &self.plugin_instances
415    }
416
417    /// Returns the exact Capability bindings in deterministic Plan order.
418    pub fn capability_bindings(&self) -> &[CapabilityBinding] {
419        &self.capability_bindings
420    }
421
422    /// Returns the Plan-declared Execution Lanes in deterministic identity order.
423    pub fn execution_lanes(&self) -> &[ExecutionLanePlan] {
424        &self.execution_lanes
425    }
426
427    /// Returns the bounded admission policy materialized for one binding Operation.
428    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    /// Returns the bounded volatile mailbox policy for one Event binding.
451    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    /// Returns the exact Plugin Instance selected by its App-local key.
470    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    /// Returns the restart policy materialized for one Plugin Instance.
477    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    /// Returns the criticality materialized for one Plugin Instance.
483    pub fn criticality_for(&self, instance_key: &str) -> Option<PluginCriticality> {
484        self.plugin_instance(instance_key)
485            .map(PluginInstancePlan::criticality)
486    }
487
488    /// Returns whether a Plugin Instance is directly bound to a required `one` Capability path.
489    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    /// Returns whether exhaustion of this Plugin Instance is terminal under
504    /// the Plan-selected failure policy.
505    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}