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 fn requirements() -> &'static [Requirement] {
255 &[]
256 }
257
258 fn create(context: CreateContext) -> Result<Self, String>
260 where
261 Self: Sized;
262
263 fn stop(&self, _context: Ctx) -> Result<(), String> {
265 Ok(())
266 }
267}
268
269impl<T> Plugin for T
270where
271 T: Default + Send + Sync + 'static,
272{
273 fn create(_context: CreateContext) -> Result<Self, String> {
274 Ok(Self::default())
275 }
276}
277
278#[cfg(not(target_arch = "wasm32"))]
279#[doc(hidden)]
280pub fn __process_descriptor(base: &str, requirements: &[Requirement]) -> String {
281 let mut value: Value = serde_json::from_str(base).expect("generated descriptor is valid JSON");
282 let object = value
283 .as_object_mut()
284 .expect("generated descriptor is a JSON object");
285 if requirements.is_empty() {
286 return base.to_owned();
287 }
288 object.insert(
289 "abi".to_owned(),
290 Value::String("lenso.json-host-imports@2".to_owned()),
291 );
292 let mut requirements = requirements
293 .iter()
294 .map(|requirement| {
295 serde_json::json!({
296 "requirement_id": requirement.requirement_id,
297 "capability_id": requirement.capability_id,
298 "descriptor_version": requirement.descriptor_version,
299 "cardinality": match requirement.cardinality {
300 Cardinality::One => "one",
301 Cardinality::Optional => "optional",
302 Cardinality::Many => "many",
303 },
304 })
305 })
306 .collect::<Vec<_>>();
307 requirements.sort_by(|left, right| {
308 left["requirement_id"]
309 .as_str()
310 .cmp(&right["requirement_id"].as_str())
311 });
312 object.insert(
313 "required_capabilities".to_owned(),
314 Value::Array(requirements),
315 );
316 serde_json::to_string(&value).expect("generated descriptor remains valid JSON")
317}
318
319#[cfg(not(target_arch = "wasm32"))]
320#[doc(hidden)]
321pub fn __validate_process_initialization<P: Plugin>(
322 params: &lenso_process_sdk::authoring::InitializeParams,
323) -> Result<(), String> {
324 use lenso_process_sdk::authoring::RequirementCardinality;
325
326 let declared = params
327 .required_declarations
328 .iter()
329 .map(|requirement| {
330 (
331 requirement.requirement_id.as_str(),
332 requirement.capability_id.as_str(),
333 requirement.descriptor_version.as_str(),
334 requirement.descriptor_digest.as_str(),
335 requirement.cardinality,
336 )
337 })
338 .collect::<Vec<_>>();
339 let expected = P::requirements()
340 .iter()
341 .map(|requirement| {
342 (
343 requirement.requirement_id,
344 requirement.capability_id,
345 requirement.descriptor_version,
346 requirement.descriptor_digest,
347 match requirement.cardinality {
348 Cardinality::One => RequirementCardinality::One,
349 Cardinality::Optional => RequirementCardinality::Optional,
350 Cardinality::Many => RequirementCardinality::Many,
351 },
352 )
353 })
354 .collect::<Vec<_>>();
355 if declared == expected {
356 Ok(())
357 } else {
358 Err("Host initialization does not match source-declared dependencies".to_owned())
359 }
360}
361
362#[cfg(not(target_arch = "wasm32"))]
363#[doc(hidden)]
364pub fn __process_create_context(
365 initialization: &lenso_process_sdk::authoring::InitializeParams,
366 call: lenso_process_sdk::ProcessLifecycleContext,
367) -> CreateContext {
368 let routes = initialization
369 .routes
370 .iter()
371 .map(|route| Dependency {
372 requirement_id: route.requirement_id.clone(),
373 route_id: route.route_id.clone(),
374 provider_instance: route.provider_instance.clone(),
375 capability_id: route.capability_id.clone(),
376 descriptor_version: route.descriptor_version.clone(),
377 descriptor_digest: route.descriptor_digest.clone(),
378 })
379 .collect();
380 CreateContext {
381 config: initialization.config.clone(),
382 dependencies: Dependencies { routes },
383 call: Ctx { inner: call },
384 }
385}
386
387#[cfg(not(target_arch = "wasm32"))]
388#[doc(hidden)]
389pub fn __process_ctx(call: lenso_process_sdk::ProcessCallContext) -> Ctx {
390 Ctx { inner: call }
391}
392
393#[cfg(target_arch = "wasm32")]
394#[doc(hidden)]
395pub mod __wasm {
396 wit_bindgen::generate!({
397 inline: r#"
398 package lenso:runtime@1.0.0;
399
400 world plugin {
401 export describe: func() -> string;
402 export invoke: func(capability: string, operation: string, request-json: string) -> result<string, string>;
403 }
404 "#,
405 world: "plugin",
406 export_macro_name: "export_lenso_plugin",
407 pub_export_macro: true,
408 default_bindings_module: "::lenso::__wasm",
409 });
410}
411
412#[doc(hidden)]
414#[derive(Clone, Debug)]
415pub enum InvocationOutcome {
416 Success(Value),
418 DomainError(Value),
420 Failure(String),
422}
423
424#[doc(hidden)]
429pub trait JsonRequestHandler {
430 fn invoke(&self, _capability: &str, _operation: &str, _request: Value) -> InvocationOutcome {
432 InvocationOutcome::Failure("this Plugin requires invocation context".to_owned())
433 }
434
435 fn invoke_with_context(
437 &self,
438 _context: Ctx,
439 capability: &str,
440 operation: &str,
441 request: Value,
442 ) -> InvocationOutcome {
443 self.invoke(capability, operation, request)
444 }
445}
446
447#[doc(hidden)]
451#[macro_export]
452macro_rules! __export_json_request_handler {
453 (
454 $plugin:ty {
455 capability_id: $capability_id:literal,
456 descriptor_version: $descriptor_version:literal,
457 descriptor_digest: $descriptor_digest:literal,
458 requests: [$first_request:literal $(, $request:literal)* $(,)?] $(,)?
459 }
460 ) => {
461 type __LensoExportedPlugin = $plugin;
462
463 #[cfg(target_arch = "wasm32")]
464 mod __lenso_wasm_export {
465 struct Component;
466
467 const DESCRIPTOR: &str = $crate::__private::lenso_guest_sdk::__request_plugin_descriptor!(
468 $capability_id,
469 $descriptor_version,
470 $first_request $(, $request)*
471 );
472
473 #[used]
474 #[unsafe(link_section = "lenso.plugin-descriptor.v1")]
475 static DESCRIPTOR_SECTION: [u8; DESCRIPTOR.len()] =
476 $crate::__private::lenso_guest_sdk::__descriptor_bytes(DESCRIPTOR);
477
478 impl $crate::__private::wasm::Guest for Component {
479 fn describe() -> ::std::string::String {
480 DESCRIPTOR.to_owned()
481 }
482
483 fn invoke(
484 capability: ::std::string::String,
485 operation: ::std::string::String,
486 request_json: ::std::string::String,
487 ) -> ::std::result::Result<::std::string::String, ::std::string::String> {
488 let request = $crate::__private::serde_json::from_str(&request_json)
489 .map_err(|_| "\"invalid_arguments\"".to_owned())?;
490 match <super::__LensoExportedPlugin as $crate::JsonRequestHandler>::invoke(
491 &<super::__LensoExportedPlugin as ::std::default::Default>::default(),
492 &capability,
493 &operation,
494 request,
495 ) {
496 $crate::InvocationOutcome::Success(value) =>
497 $crate::__private::serde_json::to_string(&value)
498 .map_err(|error| error.to_string()),
499 $crate::InvocationOutcome::DomainError(value) =>
500 ::std::result::Result::Err(
501 $crate::__private::serde_json::to_string(&value)
502 .unwrap_or_else(|_| "\"execution_failed\"".to_owned()),
503 ),
504 $crate::InvocationOutcome::Failure(detail) =>
505 ::std::result::Result::Err(
506 $crate::__private::serde_json::to_string(&detail)
507 .unwrap_or_else(|_| "\"execution_failed\"".to_owned()),
508 ),
509 }
510 }
511 }
512
513 $crate::__private::wasm::export_lenso_plugin!(Component);
514 }
515
516 #[cfg(not(target_arch = "wasm32"))]
517 mod __lenso_process_export {
518 #[derive(Default)]
519 struct Component {
520 initialization: ::std::sync::Mutex<::std::option::Option<
521 $crate::__private::lenso_process_sdk::authoring::InitializeParams,
522 >>,
523 }
524
525 const DESCRIPTOR: &str =
526 $crate::__private::lenso_guest_sdk::__request_plugin_descriptor!(
527 $capability_id,
528 $descriptor_version,
529 digest: $descriptor_digest,
530 $first_request $(, $request)*
531 );
532
533 impl $crate::__private::lenso_process_sdk::ProcessPluginV2 for Component {
534 type Instance = super::__LensoExportedPlugin;
535
536 fn initialize(
537 &self,
538 params: &$crate::__private::lenso_process_sdk::authoring::InitializeParams,
539 ) -> ::std::result::Result<(), ::std::string::String> {
540 $crate::__validate_process_initialization::<super::__LensoExportedPlugin>(params)?;
541 let mut initialization = self
542 .initialization
543 .lock()
544 .map_err(|_| "Plugin initialization state was poisoned".to_owned())?;
545 if initialization.replace(params.clone()).is_some() {
546 return Err("Plugin initialized more than once".to_owned());
547 }
548 Ok(())
549 }
550
551 fn construct(
552 &self,
553 _params: &$crate::__private::lenso_process_sdk::authoring::ConstructParams,
554 context: $crate::__private::lenso_process_sdk::ProcessLifecycleContext,
555 ) -> ::std::result::Result<Self::Instance, ::std::string::String> {
556 let initialization = self
557 .initialization
558 .lock()
559 .map_err(|_| "Plugin initialization state was poisoned".to_owned())?
560 .clone()
561 .ok_or_else(|| "Plugin constructed before initialization".to_owned())?;
562 <super::__LensoExportedPlugin as $crate::Plugin>::create(
563 $crate::__process_create_context(&initialization, context),
564 )
565 }
566
567 fn invoke(
568 &self,
569 instance: &Self::Instance,
570 params: $crate::__private::lenso_process_sdk::authoring::InvokeParams,
571 context: $crate::__private::lenso_process_sdk::ProcessInvocationContext,
572 ) -> $crate::__private::lenso_process_sdk::authoring::InvocationOutcome {
573 match <super::__LensoExportedPlugin as $crate::JsonRequestHandler>::invoke_with_context(
574 instance,
575 $crate::__process_ctx(context),
576 ¶ms.capability_id,
577 ¶ms.operation,
578 params.payload,
579 ) {
580 $crate::InvocationOutcome::Success(value) =>
581 $crate::__private::lenso_process_sdk::authoring::InvocationOutcome::Success {
582 value,
583 },
584 $crate::InvocationOutcome::DomainError(value) =>
585 $crate::__private::lenso_process_sdk::authoring::InvocationOutcome::Domain {
586 error: value,
587 },
588 $crate::InvocationOutcome::Failure(detail) =>
589 $crate::__private::lenso_process_sdk::authoring::InvocationOutcome::Runtime {
590 failure: $crate::__private::lenso_process_sdk::authoring::RuntimeFailure::PluginFailure {
591 detail,
592 },
593 },
594 }
595 }
596
597 fn stop(
598 &self,
599 instance: &Self::Instance,
600 _params: &$crate::__private::lenso_process_sdk::authoring::StopParams,
601 context: $crate::__private::lenso_process_sdk::ProcessLifecycleContext,
602 ) -> $crate::__private::lenso_process_sdk::ProcessStopOutcome {
603 match <super::__LensoExportedPlugin as $crate::Plugin>::stop(
604 instance,
605 $crate::__process_ctx(context),
606 ) {
607 Ok(()) => $crate::__private::lenso_process_sdk::ProcessStopOutcome::Completed,
608 Err(detail) => $crate::__private::lenso_process_sdk::ProcessStopOutcome::Failed(detail),
609 }
610 }
611 }
612
613 pub fn serve() {
614 if ::std::env::args_os().nth(1).as_deref()
615 == Some(::std::ffi::OsStr::new("--lenso-describe"))
616 {
617 println!(
618 "{}",
619 $crate::__process_descriptor(
620 DESCRIPTOR,
621 <super::__LensoExportedPlugin as $crate::Plugin>::requirements(),
622 ),
623 );
624 return;
625 }
626 $crate::__private::lenso_process_sdk::serve_v2(Component::default())
627 .expect("serve Lenso Process V2 Plugin");
628 }
629 }
630
631 #[cfg(not(target_arch = "wasm32"))]
632 fn main() {
633 __lenso_process_export::serve();
634 }
635 };
636}
637
638#[doc(hidden)]
639pub mod __private {
640 pub use serde_json;
641
642 #[cfg(target_arch = "wasm32")]
643 pub use crate::__wasm as wasm;
644 pub use lenso_guest_sdk;
645 #[cfg(not(target_arch = "wasm32"))]
646 pub use lenso_process_sdk;
647}
648
649#[cfg(all(test, not(target_arch = "wasm32")))]
650mod tests {
651 use super::*;
652 use lenso_process_sdk::authoring::{
653 AuthoringLimits, InitializeParams, ProvidedEndpoint, RequirementCardinality,
654 RequirementDeclaration, RouteDescriptor, SessionIdentity,
655 };
656
657 const STORE_DIGEST: &str =
658 "sha256:1111111111111111111111111111111111111111111111111111111111111111";
659 const REQUIREMENTS: &[Requirement] = &[
660 Requirement {
661 requirement_id: "destination",
662 capability_id: "example.store@1",
663 descriptor_version: "1.0.0",
664 descriptor_digest: STORE_DIGEST,
665 cardinality: Cardinality::One,
666 },
667 Requirement {
668 requirement_id: "source",
669 capability_id: "example.store@1",
670 descriptor_version: "1.0.0",
671 descriptor_digest: STORE_DIGEST,
672 cardinality: Cardinality::One,
673 },
674 ];
675
676 struct Stateful;
677
678 impl Plugin for Stateful {
679 fn requirements() -> &'static [Requirement] {
680 REQUIREMENTS
681 }
682
683 fn create(_context: CreateContext) -> Result<Self, String> {
684 Ok(Self)
685 }
686 }
687
688 fn initialization() -> InitializeParams {
689 InitializeParams {
690 api_version: 2,
691 identity: SessionIdentity {
692 session: "session-1".to_owned(),
693 plugin_instance: "sync".to_owned(),
694 plugin_generation: "generation-1".to_owned(),
695 artifact_digest: STORE_DIGEST.to_owned(),
696 contract_digest: STORE_DIGEST.to_owned(),
697 runtime_profile: "lenso.process-stdio@2".to_owned(),
698 value_profile: "lenso-json-value-v1".to_owned(),
699 },
700 config: serde_json::json!({ "prefix": "copied" }),
701 required_declarations: REQUIREMENTS
702 .iter()
703 .map(|requirement| RequirementDeclaration {
704 requirement_id: requirement.requirement_id.to_owned(),
705 capability_id: requirement.capability_id.to_owned(),
706 descriptor_version: requirement.descriptor_version.to_owned(),
707 descriptor_digest: requirement.descriptor_digest.to_owned(),
708 cardinality: RequirementCardinality::One,
709 })
710 .collect(),
711 routes: REQUIREMENTS
712 .iter()
713 .enumerate()
714 .map(|(index, requirement)| RouteDescriptor {
715 route_id: format!("route-{index}"),
716 requirement_id: requirement.requirement_id.to_owned(),
717 capability_id: requirement.capability_id.to_owned(),
718 descriptor_version: requirement.descriptor_version.to_owned(),
719 descriptor_digest: requirement.descriptor_digest.to_owned(),
720 provider_instance: format!("store-{index}"),
721 provider_order: 0,
722 })
723 .collect(),
724 provided_endpoints: vec![ProvidedEndpoint {
725 endpoint_id: "sync".to_owned(),
726 capability_id: "example.sync@1".to_owned(),
727 descriptor_version: "1.0.0".to_owned(),
728 descriptor_digest: STORE_DIGEST.to_owned(),
729 }],
730 limits: AuthoringLimits::defaults(),
731 }
732 }
733
734 #[test]
735 fn process_descriptor_contains_sorted_named_requirements() {
736 let base = r#"{"abi":"lenso.json-request@1","capabilities":[{"capability_id":"example.sync@1","descriptor_version":"1.0.0","request_operations":["sync"]}]}"#;
737 let descriptor: Value =
738 serde_json::from_str(&__process_descriptor(base, REQUIREMENTS)).unwrap();
739 assert_eq!(descriptor["abi"], "lenso.json-host-imports@2");
740 assert_eq!(
741 descriptor["required_capabilities"][0]["requirement_id"],
742 "destination"
743 );
744 assert_eq!(
745 descriptor["required_capabilities"][1]["requirement_id"],
746 "source"
747 );
748 }
749
750 #[test]
751 fn process_initialization_must_match_source_declarations_exactly() {
752 let initialization = initialization();
753 assert!(__validate_process_initialization::<Stateful>(&initialization).is_ok());
754
755 let mut drifted = initialization;
756 drifted.required_declarations[0].descriptor_digest = STORE_DIGEST.replace('1', "2");
757 assert_eq!(
758 __validate_process_initialization::<Stateful>(&drifted).unwrap_err(),
759 "Host initialization does not match source-declared dependencies"
760 );
761 }
762}