use std::{fmt, rc::Rc};
use futures::future::LocalBoxFuture;
use lenso_kernel::{InvocationContext, NativeRequestEndpoint, NativeRequestFuture, NativeRequestHandle, PluginDependencies, RequestCapability, RuntimeFailure};
use lenso_plugin_authoring::{BoundCapabilityClient, CapabilityClient, CapabilityClientMany, CapabilityReference};
pub const CAPABILITY_ID: &str = "lenso.configuration.source@1";
pub const DESCRIPTOR_VERSION: &str = "1.0.0";
pub const DESCRIPTOR_DIGEST: &str = "sha256:ff815d269114b4cbe4b268a68dee06979f292395deadcc2a52a561634b48ca21";
pub const PORTABLE: bool = true;
pub const CROSS_LANE_TRANSFER: bool = false;
pub const SOURCE_CAPABILITY_ID: &str = CAPABILITY_ID;
pub const SOURCE_DESCRIPTOR_VERSION: &str = DESCRIPTOR_VERSION;
pub const SOURCE_DESCRIPTOR_DIGEST: &str = DESCRIPTOR_DIGEST;
pub const SOURCE_CONTRACT: CapabilityReference<SourceClient> = CapabilityReference::new(CAPABILITY_ID, DESCRIPTOR_VERSION, DESCRIPTOR_DIGEST);
#[doc(hidden)]
#[macro_export]
macro_rules! __lenso_provided_source { () => { "{\"capability_id\":\"lenso.configuration.source@1\",\"descriptor_version\":\"1.0.0\",\"operations\":[\"fetch\"],\"operation_kinds\":{},\"default_admission\":{\"queue_capacity\":0,\"max_concurrency\":1},\"operation_admissions\":{},\"event_admission\":null,\"cross_lane_transfer\":false}" }; }
#[doc(hidden)]
#[macro_export]
macro_rules! __lenso_required_source_client {
() => { "{\"capability_id\":\"lenso.configuration.source@1\",\"descriptor_version\":\"1.0.0\",\"cardinality\":\"one\"}" };
($requirement_id:literal) => { concat!("{\"requirement_id\":", stringify!($requirement_id), ",\"capability_id\":\"lenso.configuration.source@1\",\"descriptor_version\":\"1.0.0\",\"cardinality\":\"one\"}") };
}
#[doc(hidden)]
#[macro_export]
macro_rules! __lenso_required_optional_source_client {
($requirement_id:literal) => { concat!("{\"requirement_id\":", stringify!($requirement_id), ",\"capability_id\":\"lenso.configuration.source@1\",\"descriptor_version\":\"1.0.0\",\"cardinality\":\"optional\"}") };
}
#[doc(hidden)]
#[macro_export]
macro_rules! __lenso_required_many_source_client {
() => { "{\"capability_id\":\"lenso.configuration.source@1\",\"descriptor_version\":\"1.0.0\",\"cardinality\":\"many\"}" };
($requirement_id:literal) => { concat!("{\"requirement_id\":", stringify!($requirement_id), ",\"capability_id\":\"lenso.configuration.source@1\",\"descriptor_version\":\"1.0.0\",\"cardinality\":\"many\"}") };
}
pub const FETCH_OPERATION: &str = "fetch";
pub use lenso_contract_runtime::{Uint64, UnknownDomainError};
use lenso_contract_runtime::{decode_portable_json, encode_portable_json};
#[derive(Clone, Copy, Debug, PartialEq, serde::Serialize, serde::Deserialize)]
pub struct FetchRequest {
}
#[derive(Clone, Debug, PartialEq, serde::Serialize, serde::Deserialize)]
pub struct FetchResponse {
#[serde(rename = "configurations")]
#[serde(deserialize_with = "lenso_contract_runtime::serde::deserialize_required")]
pub configurations: Vec<FetchResponseConfigurationsItem>,
#[serde(rename = "revision")]
#[serde(deserialize_with = "lenso_contract_runtime::serde::deserialize_required")]
pub revision: Uint64,
}
#[derive(Clone, Debug, PartialEq, serde::Serialize, serde::Deserialize)]
pub struct FetchResponseConfigurationsItem {
#[serde(rename = "instance_key")]
#[serde(deserialize_with = "lenso_contract_runtime::serde::deserialize_required")]
pub instance_key: String,
#[serde(rename = "plugin_id")]
#[serde(deserialize_with = "lenso_contract_runtime::serde::deserialize_required")]
pub plugin_id: String,
#[serde(rename = "toml")]
#[serde(deserialize_with = "lenso_contract_runtime::serde::deserialize_required")]
pub toml: String,
}
#[derive(Clone, Debug, PartialEq)]
pub enum FetchError {
InvalidSource,
Unknown(UnknownDomainError),
}
#[derive(Debug)]
pub struct Source;
impl RequestCapability for Source {
type Request = FetchRequest;
type Response = FetchResponse;
type DomainError = FetchError;
const ID: &'static str = CAPABILITY_ID;
const DESCRIPTOR_VERSION: &'static str = DESCRIPTOR_VERSION;
fn invoke_native(endpoint: &dyn NativeRequestEndpoint, operation: &str, request: Self::Request, context: InvocationContext) -> NativeRequestFuture<Self> {
if operation != FETCH_OPERATION {
return lenso_kernel::invoke_typed_or_erased_native_request::<Self>(endpoint, operation, request, context);
}
let Some(typed_endpoint) = endpoint
.typed_endpoint()
.and_then(|endpoint| endpoint.downcast_ref::<SourceRequestEndpoint>())
else {
return lenso_kernel::invoke_typed_or_erased_native_request::<Self>(endpoint, operation, request, context);
};
Rc::clone(&typed_endpoint.provider).fetch(context, request)
}
}
impl serde::Serialize for FetchError {
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: serde::Serializer,
{
use serde::ser::SerializeMap;
match self {
Self::InvalidSource => serializer.serialize_str("invalid_source"),
Self::Unknown(value) => {
let mut map = serializer.serialize_map(Some(1 + usize::from(value.payload.is_some()) + value.extra.len()))?;
map.serialize_entry("code", &value.code)?;
if let Some(payload) = &value.payload {
map.serialize_entry("payload", payload)?;
}
for (key, extra) in &value.extra {
map.serialize_entry(key, extra)?;
}
map.end()
},
}
}
}
impl<'de> serde::Deserialize<'de> for FetchError {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: serde::Deserializer<'de>,
{
let value = <serde_json::Value as serde::Deserialize>::deserialize(deserializer)?;
match value {
serde_json::Value::String(code) => match code.as_str() {
"invalid_source" => Ok(Self::InvalidSource),
_ => Ok(Self::Unknown(UnknownDomainError { code, payload: None, extra: std::collections::BTreeMap::new() })),
},
serde_json::Value::Object(mut object) => {
let Some(code) = object.remove("code").and_then(|value| value.as_str().map(ToOwned::to_owned)) else {
return Err(serde::de::Error::custom("Domain Error object is missing a string code"));
};
let payload = object.remove("payload");
let extra = object.into_iter().collect::<std::collections::BTreeMap<_, _>>();
Ok(Self::Unknown(UnknownDomainError { code, payload, extra }))
}
other => Err(serde::de::Error::custom(format!("Domain Error must be a string or object, got {other}"))),
}
}
}
pub fn encode_fetch_request(value: &FetchRequest) -> Result<String, serde_json::Error> { encode_portable_json(value) }
pub fn decode_fetch_request(wire: &str) -> Result<FetchRequest, serde_json::Error> { decode_portable_json(wire) }
pub fn encode_fetch_response(value: &FetchResponse) -> Result<String, serde_json::Error> { encode_portable_json(value) }
pub fn decode_fetch_response(wire: &str) -> Result<FetchResponse, serde_json::Error> { decode_portable_json(wire) }
pub fn encode_fetch_error(value: &FetchError) -> Result<String, serde_json::Error> { encode_portable_json(value) }
pub fn decode_fetch_error(wire: &str) -> Result<FetchError, serde_json::Error> { decode_portable_json(wire) }
#[doc(hidden)]
pub trait __LensoIntoSourceFetchResult {
fn __lenso_into_result(self) -> Result<Result<FetchResponse, FetchError>, RuntimeFailure>;
}
impl __LensoIntoSourceFetchResult for Result<FetchResponse, FetchError> {
fn __lenso_into_result(self) -> Result<Result<FetchResponse, FetchError>, RuntimeFailure> { Ok(self) }
}
impl __LensoIntoSourceFetchResult for Result<Result<FetchResponse, FetchError>, RuntimeFailure> {
fn __lenso_into_result(self) -> Result<Result<FetchResponse, FetchError>, RuntimeFailure> { self }
}
impl __LensoIntoSourceFetchResult for Result<FetchResponse, lenso_plugin_authoring::PluginError<FetchError, RuntimeFailure>> {
fn __lenso_into_result(self) -> Result<Result<FetchResponse, FetchError>, RuntimeFailure> {
match self {
Ok(value) => Ok(Ok(value)),
Err(lenso_plugin_authoring::PluginError::Domain(error)) => Ok(Err(error)),
Err(lenso_plugin_authoring::PluginError::Runtime(error)) => Err(error),
}
}
}
impl __LensoIntoSourceFetchResult for Result<FetchResponse, SourceInvocationError> {
fn __lenso_into_result(self) -> Result<Result<FetchResponse, FetchError>, RuntimeFailure> {
match self {
Ok(value) => Ok(Ok(value)),
Err(SourceInvocationError::Domain(error)) => Ok(Err(error)),
Err(SourceInvocationError::Runtime(error)) => Err(error),
}
}
}
pub trait SourceProvider: fmt::Debug + 'static {
fn fetch(&self, context: InvocationContext, request: FetchRequest) -> NativeRequestFuture<Source>;
}
#[doc(hidden)]
#[macro_export]
macro_rules! __lenso_native_lower_source {
($plugin:ty, $support:path) => {
use $support as __LensoNativeSupportSource;
impl $crate::SourceProvider for $plugin {
fn fetch(&self, context: __LensoNativeSupportSource::InvocationContext, request: $crate::FetchRequest) -> __LensoNativeSupportSource::NativeRequestFuture<$crate::Source> {
let plugin = self.clone();
::std::boxed::Box::pin(async move {
let result = <$plugin>::fetch(&plugin, context, request).await;
$crate::__LensoIntoSourceFetchResult::__lenso_into_result(result)
})
}
}
};
}
#[doc(hidden)]
#[macro_export]
macro_rules! __lenso_native_lower_object_source {
($object:ty, $plugin:ty, $support:path) => {
use $support as __LensoNativeSupportSource;
impl $crate::SourceProvider for $object {
fn fetch(&self, context: __LensoNativeSupportSource::InvocationContext, request: $crate::FetchRequest) -> __LensoNativeSupportSource::NativeRequestFuture<$crate::Source> {
let object = self.clone();
::std::boxed::Box::pin(async move {
let plugin = object.get()?;
let result = <$plugin>::fetch(plugin.as_ref(), context, request).await;
$crate::__LensoIntoSourceFetchResult::__lenso_into_result(result)
})
}
}
};
}
#[doc(hidden)]
#[macro_export]
macro_rules! __lenso_native_lower_trait_object_source {
($object:ty, $plugin:ty, $support:path) => {
use $support as __LensoNativeSupportSource;
impl $crate::SourceProvider for $object {
fn fetch(&self, context: __LensoNativeSupportSource::InvocationContext, request: $crate::FetchRequest) -> __LensoNativeSupportSource::NativeRequestFuture<$crate::Source> {
let object = self.clone();
::std::boxed::Box::pin(async move {
let plugin = object.get()?;
<$plugin as $crate::SourceProvider>::fetch(plugin.as_ref(), context, request).await
})
}
}
};
}
#[derive(Debug)]
struct SourceRequestEndpoint { provider: Rc<dyn SourceProvider> }
#[derive(Debug)]
pub struct SourceEndpoint<P: SourceProvider> { provider: Rc<P>, request_endpoint: SourceRequestEndpoint }
impl<P: SourceProvider> SourceEndpoint<P> {
pub fn new(provider: P) -> Self {
let provider = Rc::new(provider);
let request_provider: Rc<dyn SourceProvider> = provider.clone();
Self { provider, request_endpoint: SourceRequestEndpoint { provider: request_provider } }
}
}
impl<P: SourceProvider> NativeRequestEndpoint for SourceEndpoint<P> {
fn capability_id(&self) -> &'static str { CAPABILITY_ID }
fn descriptor_version(&self) -> &'static str { DESCRIPTOR_VERSION }
fn operations(&self) -> &'static [&'static str] { &[
FETCH_OPERATION,
] }
fn typed_endpoint(&self) -> Option<&dyn std::any::Any> { Some(&self.request_endpoint) }
fn invoke(&self, operation: &str, request: Box<dyn std::any::Any>, context: InvocationContext) -> LocalBoxFuture<'static, Result<Result<Box<dyn std::any::Any>, Box<dyn std::any::Any>>, RuntimeFailure>> {
match operation {
FETCH_OPERATION => {
let Ok(request) = request.downcast::<FetchRequest>() else {
return Box::pin(futures::future::ready(Err(RuntimeFailure::ProtocolViolation { capability: CAPABILITY_ID })));
};
let invocation = Rc::clone(&self.provider).fetch(context, *request);
Box::pin(async move {
invocation.await.map(|result| {
result
.map(|value| Box::new(value) as Box<dyn std::any::Any>)
.map_err(|error| Box::new(error) as Box<dyn std::any::Any>)
})
})
}
_ => Box::pin(futures::future::ready(Err(RuntimeFailure::UnknownOperation { capability: CAPABILITY_ID, operation: operation.to_owned() }))),
}
}
}
#[doc(hidden)]
#[macro_export]
macro_rules! __lenso_native_endpoints_source {
($provider:expr, $support:path) => {{
use $support as __LensoNativeSupport;
let endpoint = ::std::rc::Rc::new($crate::SourceEndpoint::new($provider));
(
vec![endpoint.clone() as ::std::rc::Rc<dyn __LensoNativeSupport::NativeRequestEndpoint>],
vec![],
vec![],
)
}};
}
#[doc(hidden)]
#[macro_export]
macro_rules! __lenso_native_provide_source {
($provider:expr, $lifecycle:expr, $support:path) => {{
use $support as __LensoNativeSupport;
let (request_endpoints, stream_endpoints, event_endpoints) =
$crate::__lenso_native_endpoints_source!($provider, $support);
__LensoNativeSupport::NativePluginInstance::with_all_endpoints(
request_endpoints,
stream_endpoints,
event_endpoints,
$lifecycle,
)
}};
}
#[derive(Clone, Debug)]
pub struct SourceClient {
fetch: NativeRequestHandle<Source>,
}
impl SourceClient {
pub fn new(handle: NativeRequestHandle<Source>) -> Self {
Self { fetch: handle }
}
pub fn from_dependencies(dependencies: &PluginDependencies) -> Result<Self, RuntimeFailure> {
<Self as CapabilityClient>::from_dependencies(dependencies)
}
pub fn from_requirement(
dependencies: &PluginDependencies,
requirement_id: &str,
) -> Result<Self, RuntimeFailure> {
<Self as CapabilityClient>::from_requirement(dependencies, requirement_id)
}
pub async fn fetch(&self, request: FetchRequest) -> Result<FetchResponse, SourceInvocationError> {
self.fetch.invoke(FETCH_OPERATION, request).await
.map_err(SourceInvocationError::Runtime)?
.map_err(SourceInvocationError::Domain)
}
pub async fn fetch_with_context(&self, context: InvocationContext, request: FetchRequest) -> Result<FetchResponse, SourceInvocationError> {
self.fetch.invoke_with_context(FETCH_OPERATION, context, request).await
.map_err(SourceInvocationError::Runtime)?
.map_err(SourceInvocationError::Domain)
}
}
impl CapabilityClient for SourceClient {
type Dependencies = PluginDependencies;
type Error = RuntimeFailure;
const CAPABILITY_ID: &'static str = CAPABILITY_ID;
const DESCRIPTOR_VERSION: &'static str = DESCRIPTOR_VERSION;
fn from_dependencies(dependencies: &PluginDependencies) -> Result<Self, RuntimeFailure> {
Ok(Self {
fetch: dependencies.one::<Source>()?,
})
}
fn from_requirement(
dependencies: &PluginDependencies,
requirement_id: &str,
) -> Result<Self, RuntimeFailure> {
let dependencies = dependencies.requirement(requirement_id)?;
Self::from_dependencies(&dependencies)
}
fn already_connected() -> RuntimeFailure {
RuntimeFailure::PluginFailure {
detail: format!("Capability Port {CAPABILITY_ID} was connected more than once"),
}
}
}
impl CapabilityClientMany for SourceClient {
fn many_from_dependencies(
dependencies: &PluginDependencies,
) -> Result<Vec<BoundCapabilityClient<Self>>, RuntimeFailure> {
dependencies
.bindings()
.iter()
.filter(|binding| binding.capability_id() == CAPABILITY_ID)
.map(|binding| {
Ok(BoundCapabilityClient::new(
binding.provider_instance(),
Self {
fetch: binding.handle().ok_or(RuntimeFailure::Unavailable { capability: CAPABILITY_ID })?.typed::<Source>()?,
},
))
})
.collect()
}
fn many_from_requirement(
dependencies: &PluginDependencies,
requirement_id: &str,
) -> Result<Vec<BoundCapabilityClient<Self>>, RuntimeFailure> {
let dependencies = dependencies.requirement(requirement_id)?;
Self::many_from_dependencies(&dependencies)
}
}
#[derive(Clone, Debug, PartialEq)]
pub enum SourceInvocationError {
Domain(FetchError),
Runtime(RuntimeFailure),
}
#[derive(Clone, Copy, Debug)]
pub struct SourceGuestClient<'a, H: lenso_guest_sdk::HostImports> {
capability: lenso_guest_sdk::GuestCapability<'a, H>,
}
impl<'a, H: lenso_guest_sdk::HostImports> SourceGuestClient<'a, H> {
pub fn from_context(context: &'a lenso_guest_sdk::GuestContext<H>) -> Result<Self, lenso_guest_sdk::GuestError<serde_json::Value>> {
context
.require(CAPABILITY_ID, DESCRIPTOR_VERSION, &[FETCH_OPERATION], &[], &[])
.map(|capability| Self { capability })
}
pub fn fetch(&self, request: &FetchRequest) -> Result<FetchResponse, lenso_guest_sdk::GuestError<FetchError>> {
self.capability.request(FETCH_OPERATION, request)
}
}
#[derive(Debug, Default)]
pub struct SourceJsonCodec;
impl lenso_runtime_codec::JsonCapabilityCodec for SourceJsonCodec {
fn capability_id(&self) -> &'static str { CAPABILITY_ID }
fn descriptor_version(&self) -> &'static str { DESCRIPTOR_VERSION }
fn descriptor_digest(&self) -> &'static str { DESCRIPTOR_DIGEST }
fn request_operations(&self) -> &'static [&'static str] { &[FETCH_OPERATION] }
fn stream_operations(&self) -> &'static [&'static str] { &[] }
fn event_operations(&self) -> &'static [&'static str] { &[] }
fn encode_request(&self, operation: &str, request: &dyn std::any::Any) -> Result<serde_json::Value, RuntimeFailure> {
match operation {
FETCH_OPERATION => {
let value = request.downcast_ref::<FetchRequest>().ok_or_else(runtime_codec_protocol_failure)?;
serde_json::to_value(value).map_err(|_| runtime_codec_protocol_failure())
},
_ => Err(runtime_codec_unknown_operation(operation)),
}
}
fn decode_response(&self, operation: &str, value: serde_json::Value) -> Result<Box<dyn std::any::Any>, RuntimeFailure> {
match operation {
FETCH_OPERATION => serde_json::from_value::<FetchResponse>(value)
.map(|value| Box::new(value) as Box<dyn std::any::Any>)
.map_err(|_| runtime_codec_protocol_failure()),
_ => Err(runtime_codec_unknown_operation(operation)),
}
}
fn decode_domain_error(&self, operation: &str, value: serde_json::Value) -> Result<Box<dyn std::any::Any>, RuntimeFailure> {
match operation {
FETCH_OPERATION => serde_json::from_value::<FetchError>(value)
.map(|value| Box::new(value) as Box<dyn std::any::Any>)
.map_err(|_| runtime_codec_protocol_failure()),
_ => Err(runtime_codec_unknown_operation(operation)),
}
}
fn encode_stream_open(&self, operation: &str, _request: &dyn std::any::Any) -> Result<serde_json::Value, RuntimeFailure> {
Err(runtime_codec_unknown_operation(operation))
}
fn encode_stream_message(&self, operation: &str, _message: &dyn std::any::Any) -> Result<serde_json::Value, RuntimeFailure> {
Err(runtime_codec_unknown_operation(operation))
}
fn decode_stream_message(&self, operation: &str, _value: serde_json::Value) -> Result<Box<dyn std::any::Any>, RuntimeFailure> {
Err(runtime_codec_unknown_operation(operation))
}
fn decode_stream_domain_error(&self, operation: &str, _value: serde_json::Value) -> Result<Box<dyn std::any::Any>, RuntimeFailure> {
Err(runtime_codec_unknown_operation(operation))
}
fn encode_event(&self, operation: &str, _event: &dyn std::any::Any) -> Result<serde_json::Value, RuntimeFailure> {
Err(runtime_codec_unknown_operation(operation))
}
fn invoke_host_request(&self, dependency: lenso_kernel::PluginDependencyHandle, operation: String, request: serde_json::Value, context: InvocationContext) -> lenso_runtime_codec::JsonHostRequestFuture {
match operation.as_str() {
FETCH_OPERATION => {
let request = serde_json::from_value::<FetchRequest>(request).map_err(|_| runtime_codec_protocol_failure());
Box::pin(async move {
let request = request?;
let handle = dependency.typed::<Source>()?;
match handle.invoke_with_context(FETCH_OPERATION, context, request).await? {
Ok(response) => serde_json::to_value(response)
.map(lenso_runtime_codec::JsonInvocationOutcome::Success)
.map_err(|_| runtime_codec_protocol_failure()),
Err(error) => serde_json::to_value(error)
.map(lenso_runtime_codec::JsonInvocationOutcome::DomainError)
.map_err(|_| runtime_codec_protocol_failure()),
}
})
},
_ => Box::pin(std::future::ready(Err(runtime_codec_unknown_operation(&operation)))),
}
}
fn open_host_stream(&self, _dependency: lenso_kernel::PluginStreamDependencyHandle, operation: String, _request: serde_json::Value, _context: InvocationContext) -> lenso_runtime_codec::JsonHostStreamOpenFuture {
Box::pin(std::future::ready(Err(runtime_codec_unknown_operation(&operation))))
}
fn publish_host_event(&self, _dependency: lenso_kernel::PluginEventDependencyHandle, operation: String, _event: serde_json::Value, _context: InvocationContext) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>> {
Box::pin(std::future::ready(Err(runtime_codec_unknown_operation(&operation))))
}
}
fn runtime_codec_protocol_failure() -> RuntimeFailure { RuntimeFailure::ProtocolViolation { capability: CAPABILITY_ID } }
fn runtime_codec_unknown_operation(operation: &str) -> RuntimeFailure {
RuntimeFailure::UnknownOperation { capability: CAPABILITY_ID, operation: operation.to_owned() }
}