1use serde_json::Value;
9
10pub use lenso_plugin_sdk_macros::plugin;
11
12#[derive(Clone, Copy, Debug, Eq, PartialEq)]
14pub struct Requirement {
15 pub requirement_id: &'static str,
16 pub capability_id: &'static str,
17 pub descriptor_version: &'static str,
18 pub descriptor_digest: &'static str,
19 pub cardinality: Cardinality,
20}
21
22impl Requirement {
23 pub const fn one<C: DependencyClient>(requirement_id: &'static str) -> Self {
25 Self::new::<C>(requirement_id, Cardinality::One)
26 }
27
28 pub const fn optional<C: DependencyClient>(requirement_id: &'static str) -> Self {
30 Self::new::<C>(requirement_id, Cardinality::Optional)
31 }
32
33 pub const fn many<C: DependencyClient>(requirement_id: &'static str) -> Self {
35 Self::new::<C>(requirement_id, Cardinality::Many)
36 }
37
38 const fn new<C: DependencyClient>(
39 requirement_id: &'static str,
40 cardinality: Cardinality,
41 ) -> Self {
42 Self {
43 requirement_id,
44 capability_id: C::CAPABILITY_ID,
45 descriptor_version: C::DESCRIPTOR_VERSION,
46 descriptor_digest: C::DESCRIPTOR_DIGEST,
47 cardinality,
48 }
49 }
50}
51
52#[derive(Clone, Copy, Debug, Eq, PartialEq)]
54pub enum Cardinality {
55 One,
56 Optional,
57 Many,
58}
59
60#[derive(Clone, Debug, Eq, PartialEq)]
62pub struct Dependency {
63 requirement_id: String,
64 route_id: String,
65 provider_instance: String,
66 capability_id: String,
67 descriptor_version: String,
68 descriptor_digest: String,
69}
70
71impl Dependency {
72 pub fn requirement_id(&self) -> &str {
74 &self.requirement_id
75 }
76
77 pub fn provider_instance(&self) -> &str {
79 &self.provider_instance
80 }
81
82 pub fn client<C: DependencyClient>(self) -> Result<C, String> {
84 if self.capability_id != C::CAPABILITY_ID
85 || self.descriptor_version != C::DESCRIPTOR_VERSION
86 || self.descriptor_digest != C::DESCRIPTOR_DIGEST
87 {
88 return Err(format!(
89 "dependency `{}` does not match generated client `{}`",
90 self.requirement_id,
91 C::CAPABILITY_ID
92 ));
93 }
94 Ok(C::from_dependency(self))
95 }
96}
97
98pub trait DependencyClient: Sized {
100 const CAPABILITY_ID: &'static str;
101 const DESCRIPTOR_VERSION: &'static str;
102 const DESCRIPTOR_DIGEST: &'static str;
103
104 #[doc(hidden)]
105 fn from_dependency(dependency: Dependency) -> Self;
106}
107
108#[derive(Clone, Debug)]
110pub struct CreateContext {
111 config: Value,
112 dependencies: Dependencies,
113 call: Ctx,
114}
115
116impl CreateContext {
117 pub const fn config(&self) -> &Value {
119 &self.config
120 }
121
122 pub const fn dependencies(&self) -> &Dependencies {
124 &self.dependencies
125 }
126
127 pub const fn ctx(&self) -> &Ctx {
129 &self.call
130 }
131}
132
133#[derive(Clone, Debug, Default)]
135pub struct Dependencies {
136 routes: Vec<Dependency>,
137}
138
139impl Dependencies {
140 pub fn one(&self, requirement_id: &str) -> Result<Dependency, String> {
142 let matches = self
143 .routes
144 .iter()
145 .filter(|route| route.requirement_id == requirement_id)
146 .collect::<Vec<_>>();
147 match matches.as_slice() {
148 [route] => Ok((*route).clone()),
149 [] => Err(format!("dependency `{requirement_id}` has no provider")),
150 _ => Err(format!(
151 "dependency `{requirement_id}` has multiple providers"
152 )),
153 }
154 }
155
156 pub fn optional(&self, requirement_id: &str) -> Result<Option<Dependency>, String> {
158 let matches = self
159 .routes
160 .iter()
161 .filter(|route| route.requirement_id == requirement_id)
162 .collect::<Vec<_>>();
163 match matches.as_slice() {
164 [] => Ok(None),
165 [route] => Ok(Some((*route).clone())),
166 _ => Err(format!(
167 "dependency `{requirement_id}` has multiple providers"
168 )),
169 }
170 }
171
172 pub fn many(&self, requirement_id: &str) -> Vec<Dependency> {
174 self.routes
175 .iter()
176 .filter(|route| route.requirement_id == requirement_id)
177 .cloned()
178 .collect()
179 }
180}
181
182#[derive(Clone, Debug, PartialEq)]
184pub enum CallError<E> {
185 Domain(E),
186 Runtime(Value),
187 InvalidValue,
188}
189
190#[derive(Clone, Debug)]
192pub struct Ctx {
193 #[cfg(not(target_arch = "wasm32"))]
194 inner: lenso_process_sdk::ProcessCallContext,
195}
196
197impl Ctx {
198 #[must_use]
200 pub fn is_cancelled(&self) -> bool {
201 #[cfg(not(target_arch = "wasm32"))]
202 {
203 self.inner.is_cancelled()
204 }
205 #[cfg(target_arch = "wasm32")]
206 {
207 false
208 }
209 }
210
211 #[cfg(not(target_arch = "wasm32"))]
213 pub fn request<Request, Response, DomainError>(
214 &self,
215 dependency: &Dependency,
216 operation: &str,
217 request: &Request,
218 ) -> Result<Response, CallError<DomainError>>
219 where
220 Request: serde::Serialize,
221 Response: serde::de::DeserializeOwned,
222 DomainError: serde::de::DeserializeOwned,
223 {
224 let payload = serde_json::to_value(request).map_err(|_| CallError::InvalidValue)?;
225 match self
226 .inner
227 .call(
228 dependency.requirement_id.clone(),
229 dependency.route_id.clone(),
230 operation,
231 payload,
232 )
233 .map_err(|detail| CallError::Runtime(serde_json::json!({ "detail": detail })))?
234 {
235 lenso_process_sdk::authoring::InvocationOutcome::Success { value } => {
236 serde_json::from_value(value).map_err(|_| CallError::InvalidValue)
237 }
238 lenso_process_sdk::authoring::InvocationOutcome::Domain { error } => {
239 match serde_json::from_value(error) {
240 Ok(error) => Err(CallError::Domain(error)),
241 Err(_) => Err(CallError::InvalidValue),
242 }
243 }
244 lenso_process_sdk::authoring::InvocationOutcome::Runtime { failure } => Err(
245 CallError::Runtime(serde_json::to_value(failure).unwrap_or(Value::Null)),
246 ),
247 }
248 }
249}
250
251pub trait Plugin: Send + Sync + 'static {
253 const CONFIGURATION_SCHEMA: Option<&'static str> = None;
255
256 fn requirements() -> &'static [Requirement] {
258 &[]
259 }
260
261 fn create(context: CreateContext) -> Result<Self, String>
263 where
264 Self: Sized;
265
266 fn stop(&self, _context: Ctx) -> Result<(), String> {
268 Ok(())
269 }
270}
271
272impl<T> Plugin for T
273where
274 T: Default + Send + Sync + 'static,
275{
276 fn create(_context: CreateContext) -> Result<Self, String> {
277 Ok(Self::default())
278 }
279}
280
281#[cfg(not(target_arch = "wasm32"))]
282#[doc(hidden)]
283pub fn __process_descriptor(
284 base: &str,
285 requirements: &[Requirement],
286 configuration_schema: Option<&str>,
287) -> String {
288 let mut value: Value = serde_json::from_str(base).expect("generated descriptor is valid JSON");
289 let object = value
290 .as_object_mut()
291 .expect("generated descriptor is a JSON object");
292 if let Some(schema) = configuration_schema {
293 let schema =
294 serde_json::from_str(schema).expect("Plugin configuration schema is valid JSON");
295 object.insert("configuration_schema".to_owned(), schema);
296 }
297 if !requirements.is_empty() {
298 object.insert(
299 "abi".to_owned(),
300 Value::String("lenso.json-host-imports@2".to_owned()),
301 );
302 let mut requirements = requirements
303 .iter()
304 .map(|requirement| {
305 serde_json::json!({
306 "requirement_id": requirement.requirement_id,
307 "capability_id": requirement.capability_id,
308 "descriptor_version": requirement.descriptor_version,
309 "cardinality": match requirement.cardinality {
310 Cardinality::One => "one",
311 Cardinality::Optional => "optional",
312 Cardinality::Many => "many",
313 },
314 })
315 })
316 .collect::<Vec<_>>();
317 requirements.sort_by(|left, right| {
318 left["requirement_id"]
319 .as_str()
320 .cmp(&right["requirement_id"].as_str())
321 });
322 object.insert(
323 "required_capabilities".to_owned(),
324 Value::Array(requirements),
325 );
326 }
327 serde_json::to_string(&value).expect("generated descriptor remains valid JSON")
328}
329
330#[cfg(not(target_arch = "wasm32"))]
331#[doc(hidden)]
332pub fn __validate_process_initialization<P: Plugin>(
333 params: &lenso_process_sdk::authoring::InitializeParams,
334) -> Result<(), String> {
335 use lenso_process_sdk::authoring::RequirementCardinality;
336
337 let declared = params
338 .required_declarations
339 .iter()
340 .map(|requirement| {
341 (
342 requirement.requirement_id.as_str(),
343 requirement.capability_id.as_str(),
344 requirement.descriptor_version.as_str(),
345 requirement.descriptor_digest.as_str(),
346 requirement.cardinality,
347 )
348 })
349 .collect::<Vec<_>>();
350 let expected = P::requirements()
351 .iter()
352 .map(|requirement| {
353 (
354 requirement.requirement_id,
355 requirement.capability_id,
356 requirement.descriptor_version,
357 requirement.descriptor_digest,
358 match requirement.cardinality {
359 Cardinality::One => RequirementCardinality::One,
360 Cardinality::Optional => RequirementCardinality::Optional,
361 Cardinality::Many => RequirementCardinality::Many,
362 },
363 )
364 })
365 .collect::<Vec<_>>();
366 if declared == expected {
367 Ok(())
368 } else {
369 Err("Host initialization does not match source-declared dependencies".to_owned())
370 }
371}
372
373#[cfg(not(target_arch = "wasm32"))]
374#[doc(hidden)]
375pub fn __process_create_context(
376 initialization: &lenso_process_sdk::authoring::InitializeParams,
377 call: lenso_process_sdk::ProcessLifecycleContext,
378) -> CreateContext {
379 let routes = initialization
380 .routes
381 .iter()
382 .map(|route| Dependency {
383 requirement_id: route.requirement_id.clone(),
384 route_id: route.route_id.clone(),
385 provider_instance: route.provider_instance.clone(),
386 capability_id: route.capability_id.clone(),
387 descriptor_version: route.descriptor_version.clone(),
388 descriptor_digest: route.descriptor_digest.clone(),
389 })
390 .collect();
391 CreateContext {
392 config: initialization.config.clone(),
393 dependencies: Dependencies { routes },
394 call: Ctx { inner: call },
395 }
396}
397
398#[cfg(not(target_arch = "wasm32"))]
399#[doc(hidden)]
400pub fn __process_ctx(call: lenso_process_sdk::ProcessCallContext) -> Ctx {
401 Ctx { inner: call }
402}
403
404#[cfg(target_arch = "wasm32")]
405#[doc(hidden)]
406pub mod __wasm {
407 wit_bindgen::generate!({
408 inline: r#"
409 package lenso:runtime@1.0.0;
410
411 world plugin {
412 export describe: func() -> string;
413 export invoke: func(capability: string, operation: string, request-json: string) -> result<string, string>;
414 }
415 "#,
416 world: "plugin",
417 export_macro_name: "export_lenso_plugin",
418 pub_export_macro: true,
419 default_bindings_module: "::lenso::__wasm",
420 });
421}
422
423#[doc(hidden)]
425#[derive(Clone, Debug)]
426pub enum InvocationOutcome {
427 Success(Value),
429 DomainError(Value),
431 Failure(String),
433}
434
435#[doc(hidden)]
440pub trait JsonRequestHandler {
441 fn invoke(&self, _capability: &str, _operation: &str, _request: Value) -> InvocationOutcome {
443 InvocationOutcome::Failure("this Plugin requires invocation context".to_owned())
444 }
445
446 fn invoke_with_context(
448 &self,
449 _context: Ctx,
450 capability: &str,
451 operation: &str,
452 request: Value,
453 ) -> InvocationOutcome {
454 self.invoke(capability, operation, request)
455 }
456}
457
458#[doc(hidden)]
462#[macro_export]
463macro_rules! __export_json_request_handler {
464 (
465 $plugin:ty {
466 capability_id: $capability_id:literal,
467 descriptor_version: $descriptor_version:literal,
468 descriptor_digest: $descriptor_digest:literal,
469 requests: [$first_request:literal $(, $request:literal)* $(,)?] $(,)?
470 }
471 ) => {
472 type __LensoExportedPlugin = $plugin;
473
474 #[cfg(target_arch = "wasm32")]
475 mod __lenso_wasm_export {
476 struct Component;
477
478 const DESCRIPTOR: &str = $crate::__private::lenso_guest_sdk::__request_plugin_descriptor!(
479 $capability_id,
480 $descriptor_version,
481 digest: $descriptor_digest,
482 $first_request $(, $request)*
483 );
484
485 #[used]
486 #[unsafe(link_section = "lenso.plugin-descriptor.v1")]
487 static DESCRIPTOR_SECTION: [u8; DESCRIPTOR.len()] =
488 $crate::__private::lenso_guest_sdk::__descriptor_bytes(DESCRIPTOR);
489
490 impl $crate::__private::wasm::Guest for Component {
491 fn describe() -> ::std::string::String {
492 DESCRIPTOR.to_owned()
493 }
494
495 fn invoke(
496 capability: ::std::string::String,
497 operation: ::std::string::String,
498 request_json: ::std::string::String,
499 ) -> ::std::result::Result<::std::string::String, ::std::string::String> {
500 let request = $crate::__private::serde_json::from_str(&request_json)
501 .map_err(|_| "\"invalid_arguments\"".to_owned())?;
502 match <super::__LensoExportedPlugin as $crate::JsonRequestHandler>::invoke(
503 &<super::__LensoExportedPlugin as ::std::default::Default>::default(),
504 &capability,
505 &operation,
506 request,
507 ) {
508 $crate::InvocationOutcome::Success(value) =>
509 $crate::__private::serde_json::to_string(&value)
510 .map_err(|error| error.to_string()),
511 $crate::InvocationOutcome::DomainError(value) =>
512 ::std::result::Result::Err(
513 $crate::__private::serde_json::to_string(&value)
514 .unwrap_or_else(|_| "\"execution_failed\"".to_owned()),
515 ),
516 $crate::InvocationOutcome::Failure(detail) =>
517 ::std::result::Result::Err(
518 $crate::__private::serde_json::to_string(&detail)
519 .unwrap_or_else(|_| "\"execution_failed\"".to_owned()),
520 ),
521 }
522 }
523 }
524
525 $crate::__private::wasm::export_lenso_plugin!(Component);
526 }
527
528 #[cfg(not(target_arch = "wasm32"))]
529 mod __lenso_process_export {
530 #[derive(Default)]
531 struct Component {
532 initialization: ::std::sync::Mutex<::std::option::Option<
533 $crate::__private::lenso_process_sdk::authoring::InitializeParams,
534 >>,
535 }
536
537 const DESCRIPTOR: &str =
538 $crate::__private::lenso_guest_sdk::__request_plugin_descriptor!(
539 $capability_id,
540 $descriptor_version,
541 digest: $descriptor_digest,
542 $first_request $(, $request)*
543 );
544
545 impl $crate::__private::lenso_process_sdk::ProcessPluginV2 for Component {
546 type Instance = super::__LensoExportedPlugin;
547
548 fn initialize(
549 &self,
550 params: &$crate::__private::lenso_process_sdk::authoring::InitializeParams,
551 ) -> ::std::result::Result<(), ::std::string::String> {
552 $crate::__validate_process_initialization::<super::__LensoExportedPlugin>(params)?;
553 let mut initialization = self
554 .initialization
555 .lock()
556 .map_err(|_| "Plugin initialization state was poisoned".to_owned())?;
557 if initialization.replace(params.clone()).is_some() {
558 return Err("Plugin initialized more than once".to_owned());
559 }
560 Ok(())
561 }
562
563 fn construct(
564 &self,
565 _params: &$crate::__private::lenso_process_sdk::authoring::ConstructParams,
566 context: $crate::__private::lenso_process_sdk::ProcessLifecycleContext,
567 ) -> ::std::result::Result<Self::Instance, ::std::string::String> {
568 let initialization = self
569 .initialization
570 .lock()
571 .map_err(|_| "Plugin initialization state was poisoned".to_owned())?
572 .clone()
573 .ok_or_else(|| "Plugin constructed before initialization".to_owned())?;
574 <super::__LensoExportedPlugin as $crate::Plugin>::create(
575 $crate::__process_create_context(&initialization, context),
576 )
577 }
578
579 fn invoke(
580 &self,
581 instance: &Self::Instance,
582 params: $crate::__private::lenso_process_sdk::authoring::InvokeParams,
583 context: $crate::__private::lenso_process_sdk::ProcessInvocationContext,
584 ) -> $crate::__private::lenso_process_sdk::authoring::InvocationOutcome {
585 match <super::__LensoExportedPlugin as $crate::JsonRequestHandler>::invoke_with_context(
586 instance,
587 $crate::__process_ctx(context),
588 ¶ms.capability_id,
589 ¶ms.operation,
590 params.payload,
591 ) {
592 $crate::InvocationOutcome::Success(value) =>
593 $crate::__private::lenso_process_sdk::authoring::InvocationOutcome::Success {
594 value,
595 },
596 $crate::InvocationOutcome::DomainError(value) =>
597 $crate::__private::lenso_process_sdk::authoring::InvocationOutcome::Domain {
598 error: value,
599 },
600 $crate::InvocationOutcome::Failure(detail) =>
601 $crate::__private::lenso_process_sdk::authoring::InvocationOutcome::Runtime {
602 failure: $crate::__private::lenso_process_sdk::authoring::RuntimeFailure::PluginFailure {
603 detail,
604 },
605 },
606 }
607 }
608
609 fn stop(
610 &self,
611 instance: &Self::Instance,
612 _params: &$crate::__private::lenso_process_sdk::authoring::StopParams,
613 context: $crate::__private::lenso_process_sdk::ProcessLifecycleContext,
614 ) -> $crate::__private::lenso_process_sdk::ProcessStopOutcome {
615 match <super::__LensoExportedPlugin as $crate::Plugin>::stop(
616 instance,
617 $crate::__process_ctx(context),
618 ) {
619 Ok(()) => $crate::__private::lenso_process_sdk::ProcessStopOutcome::Completed,
620 Err(detail) => $crate::__private::lenso_process_sdk::ProcessStopOutcome::Failed(detail),
621 }
622 }
623 }
624
625 pub fn serve() {
626 if ::std::env::args_os().nth(1).as_deref()
627 == Some(::std::ffi::OsStr::new("--lenso-describe"))
628 {
629 println!(
630 "{}",
631 $crate::__process_descriptor(
632 DESCRIPTOR,
633 <super::__LensoExportedPlugin as $crate::Plugin>::requirements(),
634 <super::__LensoExportedPlugin as $crate::Plugin>::CONFIGURATION_SCHEMA,
635 ),
636 );
637 return;
638 }
639 $crate::__private::lenso_process_sdk::serve_v2(Component::default())
640 .expect("serve Lenso Process V2 Plugin");
641 }
642 }
643
644 #[cfg(not(target_arch = "wasm32"))]
645 fn main() {
646 __lenso_process_export::serve();
647 }
648 };
649}
650
651#[doc(hidden)]
652pub mod __private {
653 pub use serde_json;
654
655 #[cfg(target_arch = "wasm32")]
656 pub use crate::__wasm as wasm;
657 pub use lenso_guest_sdk;
658 #[cfg(not(target_arch = "wasm32"))]
659 pub use lenso_process_sdk;
660}
661
662#[cfg(all(test, not(target_arch = "wasm32")))]
663mod tests {
664 use super::*;
665 use lenso_process_sdk::authoring::{
666 AuthoringLimits, InitializeParams, ProvidedEndpoint, RequirementCardinality,
667 RequirementDeclaration, RouteDescriptor, SessionIdentity,
668 };
669
670 const STORE_DIGEST: &str =
671 "sha256:1111111111111111111111111111111111111111111111111111111111111111";
672 const REQUIREMENTS: &[Requirement] = &[
673 Requirement {
674 requirement_id: "destination",
675 capability_id: "example.store@1",
676 descriptor_version: "1.0.0",
677 descriptor_digest: STORE_DIGEST,
678 cardinality: Cardinality::One,
679 },
680 Requirement {
681 requirement_id: "source",
682 capability_id: "example.store@1",
683 descriptor_version: "1.0.0",
684 descriptor_digest: STORE_DIGEST,
685 cardinality: Cardinality::One,
686 },
687 ];
688
689 struct Stateful;
690
691 impl Plugin for Stateful {
692 fn requirements() -> &'static [Requirement] {
693 REQUIREMENTS
694 }
695
696 fn create(_context: CreateContext) -> Result<Self, String> {
697 Ok(Self)
698 }
699 }
700
701 fn initialization() -> InitializeParams {
702 InitializeParams {
703 api_version: 2,
704 identity: SessionIdentity {
705 session: "session-1".to_owned(),
706 plugin_instance: "sync".to_owned(),
707 plugin_generation: "generation-1".to_owned(),
708 artifact_digest: STORE_DIGEST.to_owned(),
709 contract_digest: STORE_DIGEST.to_owned(),
710 runtime_profile: "lenso.process-stdio@2".to_owned(),
711 value_profile: "lenso-json-value-v1".to_owned(),
712 },
713 config: serde_json::json!({ "prefix": "copied" }),
714 required_declarations: REQUIREMENTS
715 .iter()
716 .map(|requirement| RequirementDeclaration {
717 requirement_id: requirement.requirement_id.to_owned(),
718 capability_id: requirement.capability_id.to_owned(),
719 descriptor_version: requirement.descriptor_version.to_owned(),
720 descriptor_digest: requirement.descriptor_digest.to_owned(),
721 cardinality: RequirementCardinality::One,
722 })
723 .collect(),
724 routes: REQUIREMENTS
725 .iter()
726 .enumerate()
727 .map(|(index, requirement)| RouteDescriptor {
728 route_id: format!("route-{index}"),
729 requirement_id: requirement.requirement_id.to_owned(),
730 capability_id: requirement.capability_id.to_owned(),
731 descriptor_version: requirement.descriptor_version.to_owned(),
732 descriptor_digest: requirement.descriptor_digest.to_owned(),
733 provider_instance: format!("store-{index}"),
734 provider_order: 0,
735 })
736 .collect(),
737 provided_endpoints: vec![ProvidedEndpoint {
738 endpoint_id: "sync".to_owned(),
739 capability_id: "example.sync@1".to_owned(),
740 descriptor_version: "1.0.0".to_owned(),
741 descriptor_digest: STORE_DIGEST.to_owned(),
742 }],
743 limits: AuthoringLimits::defaults(),
744 }
745 }
746
747 #[test]
748 fn process_descriptor_contains_sorted_named_requirements() {
749 let base = r#"{"abi":"lenso.json-request@1","capabilities":[{"capability_id":"example.sync@1","descriptor_version":"1.0.0","request_operations":["sync"]}]}"#;
750 let descriptor: Value = serde_json::from_str(&__process_descriptor(
751 base,
752 REQUIREMENTS,
753 Some(r#"{"type":"object","required":["prefix"]}"#),
754 ))
755 .unwrap();
756 assert_eq!(descriptor["abi"], "lenso.json-host-imports@2");
757 assert_eq!(
758 descriptor["required_capabilities"][0]["requirement_id"],
759 "destination"
760 );
761 assert_eq!(
762 descriptor["required_capabilities"][1]["requirement_id"],
763 "source"
764 );
765 assert_eq!(descriptor["configuration_schema"]["type"], "object");
766 }
767
768 #[test]
769 fn process_initialization_must_match_source_declarations_exactly() {
770 let initialization = initialization();
771 assert!(__validate_process_initialization::<Stateful>(&initialization).is_ok());
772
773 let mut drifted = initialization;
774 drifted.required_declarations[0].descriptor_digest = STORE_DIGEST.replace('1', "2");
775 assert_eq!(
776 __validate_process_initialization::<Stateful>(&drifted).unwrap_err(),
777 "Host initialization does not match source-declared dependencies"
778 );
779 }
780}