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