1mod authoring;
4mod managed_tasks;
5
6use std::{
7 collections::BTreeMap,
8 rc::Rc,
9 sync::{Mutex, OnceLock},
10};
11
12#[doc(hidden)]
13pub use authoring::{CompleteObjectLifecycle, ConstructionContext, LifecycleContext, PluginObject};
14#[doc(hidden)]
15pub use inventory as __inventory;
16use lenso_app_plan::{
17 ExecutionClassId, ResolvedAppPlan,
18 authoring::{HostCatalog, HostDefaultPlugin, HostPluginRelease, HostSlot, PluginDescriptor},
19};
20use lenso_kernel::{ActivateContext, DeactivateContext, PrepareContext};
21pub use lenso_kernel::{CancellationToken, RuntimeFailure};
22pub use lenso_native_adapter_macros::{PluginConfig, plugin, plugin_impl, provides};
23pub use lenso_runtime_codec::InstanceResources;
24pub use managed_tasks::{ManagedTasks, ManagedTasksError};
25
26#[allow(async_fn_in_trait)]
32pub trait Lifecycle: Clone + 'static {
33 async fn prepare(&self, _context: PrepareContext) -> Result<(), RuntimeFailure> {
34 Ok(())
35 }
36
37 async fn activate(&self, _context: ActivateContext) -> Result<(), RuntimeFailure> {
38 Ok(())
39 }
40
41 async fn deactivate(&self, _context: DeactivateContext) -> Result<(), RuntimeFailure> {
42 Ok(())
43 }
44}
45
46#[doc(hidden)]
48pub mod __private {
49 pub use crate::authoring::{ErasedConstructionFuture, LinkedPluginConstruction};
50 pub use crate::{
51 __inventory, CompleteObjectLifecycle, ConstructionContext, Lifecycle, LifecycleContext,
52 LinkedNativePluginFactory, NativePluginDefinition, NativePluginFactory,
53 NativePluginFactoryContext, NativePluginInstance, PluginObject, RuntimeFailure,
54 link_native_plugin,
55 };
56 pub use futures;
57 pub use futures::future::LocalBoxFuture;
58 pub use lenso_kernel::{
59 ActivateContext, DeactivateContext, InvocationContext, NativeEventEndpoint,
60 NativeRequestEndpoint, NativeRequestFuture, NativeStreamEndpoint, NativeStreamSession,
61 PluginFuture, PluginLifecycle, PrepareContext,
62 };
63 pub use lenso_plugin_authoring::{
64 BoundCapabilityClient, CapabilityClient, CapabilityClientMany,
65 };
66 pub use lenso_runtime_codec::InstanceResources;
67 pub use serde_json;
68}
69
70use lenso_kernel::{
71 NativeEndpointSet, NativeEventEndpoint, NativeExecutionAdapter, NativeRequestEndpoint,
72 NativeStreamEndpoint, NoopPluginLifecycle, PluginLifecycle, PreparedBinding,
73 PreparedEventBinding, PreparedNativeApp, PreparedNativePlugin, PreparedStreamBinding,
74};
75
76#[derive(Clone, Copy, Debug)]
78#[doc(hidden)]
79pub struct LinkedNativePluginFactory {
80 constructor: fn() -> Rc<dyn NativePluginFactory>,
81 descriptor: &'static str,
82}
83
84impl LinkedNativePluginFactory {
85 #[doc(hidden)]
87 pub const fn new(
88 constructor: fn() -> Rc<dyn NativePluginFactory>,
89 descriptor: &'static str,
90 ) -> Self {
91 Self {
92 constructor,
93 descriptor,
94 }
95 }
96}
97
98inventory::collect!(LinkedNativePluginFactory);
99
100fn explicitly_linked_factories() -> &'static Mutex<Vec<LinkedNativePluginFactory>> {
101 static FACTORIES: OnceLock<Mutex<Vec<LinkedNativePluginFactory>>> = OnceLock::new();
102 FACTORIES.get_or_init(|| Mutex::new(Vec::new()))
103}
104
105#[doc(hidden)]
107pub fn link_native_plugin(factory: LinkedNativePluginFactory) {
108 let mut factories = explicitly_linked_factories()
109 .lock()
110 .unwrap_or_else(std::sync::PoisonError::into_inner);
111 if !factories.iter().any(|linked| {
112 linked.descriptor == factory.descriptor
113 && std::ptr::fn_addr_eq(linked.constructor, factory.constructor)
114 }) {
115 factories.push(factory);
116 }
117}
118
119fn linked_factories() -> Vec<LinkedNativePluginFactory> {
120 let factories = inventory::iter::<LinkedNativePluginFactory>
121 .into_iter()
122 .copied()
123 .collect::<Vec<_>>();
124 let mut factories = factories;
125 factories.extend(
126 explicitly_linked_factories()
127 .lock()
128 .unwrap_or_else(std::sync::PoisonError::into_inner)
129 .iter()
130 .copied(),
131 );
132 factories
133 .into_iter()
134 .fold(Vec::new(), |mut unique, factory| {
135 if !unique.iter().any(|linked: &LinkedNativePluginFactory| {
136 linked.descriptor == factory.descriptor
137 && std::ptr::fn_addr_eq(linked.constructor, factory.constructor)
138 }) {
139 unique.push(factory);
140 }
141 unique
142 })
143}
144
145#[derive(Debug)]
147pub struct NativePluginInstance {
148 endpoints: NativeEndpointSet,
149 lifecycle: Rc<dyn PluginLifecycle>,
150}
151
152impl NativePluginInstance {
153 pub fn new(endpoints: Vec<Rc<dyn NativeRequestEndpoint>>) -> Self {
155 Self::with_lifecycle(endpoints, NoopPluginLifecycle)
156 }
157
158 pub fn with_lifecycle(
160 endpoints: Vec<Rc<dyn NativeRequestEndpoint>>,
161 lifecycle: impl PluginLifecycle,
162 ) -> Self {
163 Self {
164 endpoints: NativeEndpointSet::new(endpoints, Vec::new(), Vec::new()),
165 lifecycle: Rc::new(lifecycle),
166 }
167 }
168
169 pub fn with_endpoints(
171 endpoints: Vec<Rc<dyn NativeRequestEndpoint>>,
172 stream_endpoints: Vec<Rc<dyn NativeStreamEndpoint>>,
173 lifecycle: impl PluginLifecycle,
174 ) -> Self {
175 Self {
176 endpoints: NativeEndpointSet::new(endpoints, stream_endpoints, Vec::new()),
177 lifecycle: Rc::new(lifecycle),
178 }
179 }
180
181 pub fn with_stream_endpoints(
183 stream_endpoints: Vec<Rc<dyn NativeStreamEndpoint>>,
184 lifecycle: impl PluginLifecycle,
185 ) -> Self {
186 Self::with_endpoints(Vec::new(), stream_endpoints, lifecycle)
187 }
188
189 pub fn with_event_endpoints(
191 event_endpoints: Vec<Rc<dyn NativeEventEndpoint>>,
192 lifecycle: impl PluginLifecycle,
193 ) -> Self {
194 Self {
195 endpoints: NativeEndpointSet::new(Vec::new(), Vec::new(), event_endpoints),
196 lifecycle: Rc::new(lifecycle),
197 }
198 }
199
200 pub fn with_all_endpoints(
202 endpoints: Vec<Rc<dyn NativeRequestEndpoint>>,
203 stream_endpoints: Vec<Rc<dyn NativeStreamEndpoint>>,
204 event_endpoints: Vec<Rc<dyn NativeEventEndpoint>>,
205 lifecycle: impl PluginLifecycle,
206 ) -> Self {
207 Self {
208 endpoints: NativeEndpointSet::new(endpoints, stream_endpoints, event_endpoints),
209 lifecycle: Rc::new(lifecycle),
210 }
211 }
212
213 pub fn lifecycle(&self) -> Rc<dyn PluginLifecycle> {
215 self.lifecycle.clone()
216 }
217
218 pub fn endpoints(&self) -> &[Rc<dyn NativeRequestEndpoint>] {
220 self.endpoints.request()
221 }
222
223 pub fn stream_endpoints(&self) -> &[Rc<dyn NativeStreamEndpoint>] {
225 self.endpoints.stream()
226 }
227
228 pub fn event_endpoints(&self) -> &[Rc<dyn NativeEventEndpoint>] {
230 self.endpoints.event()
231 }
232}
233
234impl Default for NativePluginInstance {
235 fn default() -> Self {
236 Self::new(Vec::new())
237 }
238}
239
240pub trait NativePluginFactory: std::fmt::Debug + 'static {
242 fn package_id(&self) -> &'static str;
244 fn package_version(&self) -> &'static str {
246 ""
247 }
248 fn runtime_profile(&self) -> &'static str {
253 "lenso.native-authoring@1"
254 }
255 fn factory_identity(&self) -> String {
262 let version = self.package_version();
263 if version.is_empty() {
264 self.package_id().to_owned()
265 } else {
266 format!("{}@{version}", self.package_id())
267 }
268 }
269 fn instantiate(
271 &self,
272 context: NativePluginFactoryContext<'_>,
273 ) -> Result<NativePluginInstance, RuntimeFailure>;
274}
275
276pub trait NativePluginDefinition: Sized + 'static {
283 const PACKAGE_ID: &'static str;
284 const PACKAGE_VERSION: &'static str;
285 const RUNTIME_PROFILE: &'static str;
286
287 fn link();
289
290 fn instantiate_with(
291 context: NativePluginFactoryContext<'_>,
292 initialize: &dyn Fn(&mut Self) -> Result<(), RuntimeFailure>,
293 ) -> Result<NativePluginInstance, RuntimeFailure>;
294}
295
296pub fn link_plugin<P: NativePluginDefinition>() {
298 P::link();
299}
300
301pub struct ConfiguredPluginFactory<P, F> {
303 initialize: F,
304 plugin: std::marker::PhantomData<fn() -> P>,
305}
306
307impl<P, F> std::fmt::Debug for ConfiguredPluginFactory<P, F> {
308 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
309 formatter
310 .debug_struct("ConfiguredPluginFactory")
311 .finish_non_exhaustive()
312 }
313}
314
315impl<P: NativePluginDefinition, F: Fn(&mut P) -> Result<(), RuntimeFailure> + 'static>
316 ConfiguredPluginFactory<P, F>
317{
318 pub fn new(initialize: F) -> Self {
319 Self {
320 initialize,
321 plugin: std::marker::PhantomData,
322 }
323 }
324}
325
326impl<P: NativePluginDefinition, F: Fn(&mut P) -> Result<(), RuntimeFailure> + 'static>
327 NativePluginFactory for ConfiguredPluginFactory<P, F>
328{
329 fn package_id(&self) -> &'static str {
330 P::PACKAGE_ID
331 }
332 fn package_version(&self) -> &'static str {
333 P::PACKAGE_VERSION
334 }
335 fn runtime_profile(&self) -> &'static str {
336 P::RUNTIME_PROFILE
337 }
338 fn instantiate(
339 &self,
340 context: NativePluginFactoryContext<'_>,
341 ) -> Result<NativePluginInstance, RuntimeFailure> {
342 P::instantiate_with(context, &self.initialize)
343 }
344}
345
346#[derive(Clone, Copy, Debug)]
348pub struct NativePluginFactoryContext<'a> {
349 instance_key: &'a str,
350 entrypoint: &'a str,
351 configuration: &'a str,
352 resources: &'a InstanceResources,
353}
354
355impl<'a> NativePluginFactoryContext<'a> {
356 fn from_plan(
357 instance: &'a lenso_app_plan::PluginInstancePlan,
358 resources: &'a InstanceResources,
359 ) -> Self {
360 Self {
361 instance_key: instance.instance_key(),
362 entrypoint: instance.entrypoint(),
363 configuration: instance.configuration(),
364 resources,
365 }
366 }
367
368 pub const fn instance_key(self) -> &'a str {
370 self.instance_key
371 }
372
373 pub const fn entrypoint(self) -> &'a str {
375 self.entrypoint
376 }
377
378 pub const fn configuration(self) -> &'a str {
380 self.configuration
381 }
382
383 pub const fn resources(self) -> &'a InstanceResources {
385 self.resources
386 }
387}
388
389#[derive(Debug, Default)]
391pub struct NativePluginRegistry {
392 factories: Vec<Rc<dyn NativePluginFactory>>,
393 resources: lenso_runtime_codec::InstanceResourceCatalog,
394 overrides: BTreeMap<String, Rc<dyn NativePluginFactory>>,
395 linked: bool,
396}
397
398type NativeInstances = BTreeMap<String, NativePluginInstance>;
399type PreparedGenerations = BTreeMap<String, PreparedNativePlugin>;
400type NativeBindings = (
401 Vec<PreparedBinding>,
402 Vec<PreparedStreamBinding>,
403 Vec<PreparedEventBinding>,
404);
405
406fn factory_matches(
407 factory: &dyn NativePluginFactory,
408 expected: &lenso_app_plan::PluginInstancePlan,
409) -> bool {
410 factory.package_id() == expected.package_id()
411 && (expected.package_revision().is_empty()
412 || factory.package_version() == expected.package_revision()
413 || factory.factory_identity() == expected.package_revision())
414}
415
416impl NativePluginRegistry {
417 pub fn new() -> Self {
419 Self::default()
420 }
421
422 #[must_use]
427 pub fn with_linked_factories(mut self) -> Self {
428 if !self.linked {
429 self.factories.extend(
430 linked_factories()
431 .into_iter()
432 .map(|linked| (linked.constructor)()),
433 );
434 self.linked = true;
435 }
436 self
437 }
438
439 #[must_use]
441 pub fn with_resources(
442 mut self,
443 resources: lenso_runtime_codec::InstanceResourceCatalog,
444 ) -> Self {
445 self.resources = resources;
446 self
447 }
448
449 pub fn factories(&self) -> impl Iterator<Item = &dyn NativePluginFactory> {
451 self.factories
452 .iter()
453 .filter(|factory| !self.overrides.contains_key(&factory.factory_identity()))
454 .chain(self.overrides.values())
455 .map(std::convert::AsRef::as_ref)
456 }
457
458 pub fn host_catalog(
460 slots: impl IntoIterator<Item = HostSlot>,
461 defaults: impl IntoIterator<Item = HostDefaultPlugin>,
462 ) -> Result<HostCatalog, RuntimeFailure> {
463 let plugins = linked_factories()
464 .into_iter()
465 .map(|linked| {
466 serde_json::from_str::<PluginDescriptor>(linked.descriptor)
467 .map(HostPluginRelease::new)
468 .map_err(|error| RuntimeFailure::InvalidResolvedPlan {
469 detail: format!("invalid linked Plugin Descriptor: {error}"),
470 })
471 })
472 .collect::<Result<Vec<_>, _>>()?;
473 Ok(HostCatalog::new(slots, plugins, defaults))
474 }
475 #[must_use]
477 pub fn with_factory(mut self, factory: impl NativePluginFactory) -> Self {
478 self.factories.push(Rc::new(factory));
479 self
480 }
481
482 pub fn with_factory_override(
486 mut self,
487 factory: impl NativePluginFactory,
488 ) -> Result<Self, RuntimeFailure> {
489 let identity = factory.factory_identity();
490 if self.overrides.contains_key(&identity) {
491 return invalid(format!("duplicate factory override `{identity}`"));
492 }
493 self.overrides.insert(identity, Rc::new(factory));
494 Ok(self)
495 }
496
497 fn validate_overrides(&self) -> Result<(), RuntimeFailure> {
498 for (identity, replacement) in &self.overrides {
499 let originals: Vec<_> = self
500 .factories
501 .iter()
502 .filter(|factory| factory.factory_identity() == *identity)
503 .collect();
504 if originals.len() != 1 {
505 return invalid(format!(
506 "factory override `{identity}` requires exactly one original"
507 ));
508 }
509 let original = &originals[0];
510 if original.package_id() != replacement.package_id()
511 || original.package_version() != replacement.package_version()
512 || original.runtime_profile() != replacement.runtime_profile()
513 {
514 return invalid(format!(
515 "factory override `{identity}` changes package metadata"
516 ));
517 }
518 }
519 Ok(())
520 }
521
522 fn prepare_instances(
523 &self,
524 plan: &ResolvedAppPlan,
525 ) -> Result<(NativeInstances, PreparedGenerations), RuntimeFailure> {
526 self.validate_overrides()?;
527 let mut instances = BTreeMap::new();
528 let mut generations = BTreeMap::new();
529 for expected in plan
530 .plugin_instances()
531 .iter()
532 .filter(|instance| instance.execution_class() == &ExecutionClassId::native_rust())
533 {
534 let matching_factories: Vec<_> = self
535 .factories()
536 .filter(|factory| factory_matches(*factory, expected))
537 .collect();
538 let factory = match matching_factories.as_slice() {
539 [] => {
540 return Err(RuntimeFailure::MissingPluginFactory {
541 instance: expected.instance_key().to_owned(),
542 package_id: expected.package_id().to_owned(),
543 });
544 }
545 [factory] => *factory,
546 _ => {
547 return invalid(format!(
548 "multiple statically linked factories declare package `{}`",
549 expected.package_id()
550 ));
551 }
552 };
553 let generation = factory.instantiate(NativePluginFactoryContext::from_plan(
554 expected,
555 self.resources.for_instance(expected.instance_key()),
556 ))?;
557 generations.insert(
558 expected.instance_key().to_owned(),
559 PreparedNativePlugin::with_endpoint_set_lifecycle(
560 generation.endpoints.clone(),
561 generation.lifecycle(),
562 ),
563 );
564 if instances
565 .insert(expected.instance_key().to_owned(), generation)
566 .is_some()
567 {
568 return invalid(format!(
569 "duplicate Plugin Instance `{}`",
570 expected.instance_key()
571 ));
572 }
573 }
574 Ok((instances, generations))
575 }
576}
577
578impl NativeExecutionAdapter for NativePluginRegistry {
579 fn supports_runtime_profile(&self, authoring_version: u32, profile: &str) -> bool {
580 matches!(
581 (authoring_version, profile),
582 (1, "lenso.native-authoring@1" | "lenso.native-rust@1")
583 | (2, "lenso.native-authoring@2")
584 )
585 }
586
587 fn prepare(&self, plan: &ResolvedAppPlan) -> Result<PreparedNativeApp, RuntimeFailure> {
588 plan.validate()
589 .map_err(|error| RuntimeFailure::InvalidResolvedPlan {
590 detail: error.to_string(),
591 })?;
592
593 let (instances, generations) = self.prepare_instances(plan)?;
594 let (bindings, stream_bindings, event_bindings) = prepare_bindings(plan, &instances)?;
595 Ok(PreparedNativeApp::new(bindings, generations)
596 .with_stream_bindings(stream_bindings)
597 .with_event_bindings(event_bindings))
598 }
599
600 fn recreate(
601 &self,
602 plan: &ResolvedAppPlan,
603 instance_key: &str,
604 ) -> Result<PreparedNativePlugin, RuntimeFailure> {
605 self.validate_overrides()?;
606 let expected = plan
607 .plugin_instances()
608 .iter()
609 .find(|instance| instance.instance_key() == instance_key)
610 .ok_or_else(|| RuntimeFailure::InvalidResolvedPlan {
611 detail: format!("unknown Plugin Instance `{instance_key}`"),
612 })?;
613 let matching_factories: Vec<_> = self
614 .factories()
615 .filter(|factory| factory_matches(*factory, expected))
616 .collect();
617 let factory = match matching_factories.as_slice() {
618 [] => {
619 return Err(RuntimeFailure::MissingPluginFactory {
620 instance: expected.instance_key().to_owned(),
621 package_id: expected.package_id().to_owned(),
622 });
623 }
624 [factory] => *factory,
625 _ => {
626 return invalid(format!(
627 "multiple statically linked factories declare package `{}`",
628 expected.package_id()
629 ));
630 }
631 };
632 let generation = factory.instantiate(NativePluginFactoryContext::from_plan(
633 expected,
634 self.resources.for_instance(expected.instance_key()),
635 ))?;
636 Ok(PreparedNativePlugin::with_endpoint_set_lifecycle(
637 generation.endpoints.clone(),
638 generation.lifecycle(),
639 ))
640 }
641}
642
643fn prepare_bindings(
644 plan: &ResolvedAppPlan,
645 instances: &NativeInstances,
646) -> Result<NativeBindings, RuntimeFailure> {
647 let mut bindings = Vec::new();
648 let mut stream_bindings = Vec::new();
649 let mut event_bindings = Vec::new();
650 for binding in plan.capability_bindings() {
651 if !instances.contains_key(binding.provider_instance()) {
652 continue;
653 }
654 let provider = plan
655 .plugin_instance(binding.provider_instance())
656 .expect("validated binding provider should exist");
657 let descriptor = provider
658 .provided_capabilities()
659 .iter()
660 .find(|descriptor| descriptor.capability_id() == binding.capability_id())
661 .expect("validated binding descriptor should exist");
662 if !descriptor.request_operations().is_empty() {
663 let endpoint = instances
664 .get(binding.provider_instance())
665 .and_then(|instance| {
666 instance.endpoints.request().iter().find(|endpoint| {
667 endpoint.capability_id() == binding.capability_id()
668 && endpoint.descriptor_version() == binding.descriptor_version()
669 })
670 })
671 .cloned()
672 .ok_or_else(|| RuntimeFailure::InvalidResolvedPlan {
673 detail: format!(
674 "Capability `{}` version `{}` has no request endpoint on provider `{}`",
675 binding.capability_id(),
676 binding.descriptor_version(),
677 binding.provider_instance()
678 ),
679 })?;
680 bindings.push(
681 PreparedBinding::new(
682 binding.consumer_instance(),
683 binding.provider_instance(),
684 endpoint,
685 )
686 .with_requirement_id(binding.requirement_id()),
687 );
688 }
689 if !descriptor.stream_operations().is_empty() {
690 let endpoint = instances
691 .get(binding.provider_instance())
692 .and_then(|instance| {
693 instance.endpoints.stream().iter().find(|endpoint| {
694 endpoint.capability_id() == binding.capability_id()
695 && endpoint.descriptor_version() == binding.descriptor_version()
696 })
697 })
698 .cloned()
699 .ok_or_else(|| RuntimeFailure::InvalidResolvedPlan {
700 detail: format!(
701 "Capability `{}` version `{}` has no stream endpoint on provider `{}`",
702 binding.capability_id(),
703 binding.descriptor_version(),
704 binding.provider_instance()
705 ),
706 })?;
707 stream_bindings.push(
708 PreparedStreamBinding::new(
709 binding.consumer_instance(),
710 binding.provider_instance(),
711 endpoint,
712 )
713 .with_requirement_id(binding.requirement_id()),
714 );
715 }
716 if !descriptor.event_operations().is_empty() {
717 let endpoint = instances
718 .get(binding.provider_instance())
719 .and_then(|instance| {
720 instance.endpoints.event().iter().find(|endpoint| {
721 endpoint.capability_id() == binding.capability_id()
722 && endpoint.descriptor_version() == binding.descriptor_version()
723 })
724 })
725 .cloned()
726 .ok_or_else(|| RuntimeFailure::InvalidResolvedPlan {
727 detail: format!(
728 "Capability `{}` version `{}` has no Event endpoint on provider `{}`",
729 binding.capability_id(),
730 binding.descriptor_version(),
731 binding.provider_instance()
732 ),
733 })?;
734 event_bindings.push(
735 PreparedEventBinding::new(
736 binding.consumer_instance(),
737 binding.provider_instance(),
738 endpoint,
739 )
740 .with_requirement_id(binding.requirement_id()),
741 );
742 }
743 }
744 Ok((bindings, stream_bindings, event_bindings))
745}
746
747fn invalid<T>(detail: String) -> Result<T, RuntimeFailure> {
748 Err(RuntimeFailure::InvalidResolvedPlan { detail })
749}