1use std::{fmt, rc::Rc};
3use futures::future::LocalBoxFuture;
4use lenso_kernel::{InvocationContext, NativeRequestEndpoint, NativeRequestFuture, NativeRequestHandle, PluginDependencies, RequestCapability, RuntimeFailure};
5
6use lenso_plugin_authoring::{BoundCapabilityClient, CapabilityClient, CapabilityClientMany, CapabilityReference};
7pub const CAPABILITY_ID: &str = "lenso.configuration.source@1";
8pub const DESCRIPTOR_VERSION: &str = "1.0.0";
9pub const DESCRIPTOR_DIGEST: &str = "sha256:ff815d269114b4cbe4b268a68dee06979f292395deadcc2a52a561634b48ca21";
10pub const PORTABLE: bool = true;
11pub const CROSS_LANE_TRANSFER: bool = false;
12pub const SOURCE_CAPABILITY_ID: &str = CAPABILITY_ID;
13pub const SOURCE_DESCRIPTOR_VERSION: &str = DESCRIPTOR_VERSION;
14pub const SOURCE_DESCRIPTOR_DIGEST: &str = DESCRIPTOR_DIGEST;
15pub const SOURCE_CONTRACT: CapabilityReference<SourceClient> = CapabilityReference::new(CAPABILITY_ID, DESCRIPTOR_VERSION, DESCRIPTOR_DIGEST);
16
17#[doc(hidden)]
18#[macro_export]
19macro_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}" }; }
20
21#[doc(hidden)]
22#[macro_export]
23macro_rules! __lenso_required_source_client {
24 () => { "{\"capability_id\":\"lenso.configuration.source@1\",\"descriptor_version\":\"1.0.0\",\"cardinality\":\"one\"}" };
25 ($requirement_id:literal) => { concat!("{\"requirement_id\":", stringify!($requirement_id), ",\"capability_id\":\"lenso.configuration.source@1\",\"descriptor_version\":\"1.0.0\",\"cardinality\":\"one\"}") };
26}
27
28#[doc(hidden)]
29#[macro_export]
30macro_rules! __lenso_required_optional_source_client {
31 ($requirement_id:literal) => { concat!("{\"requirement_id\":", stringify!($requirement_id), ",\"capability_id\":\"lenso.configuration.source@1\",\"descriptor_version\":\"1.0.0\",\"cardinality\":\"optional\"}") };
32}
33
34#[doc(hidden)]
35#[macro_export]
36macro_rules! __lenso_required_many_source_client {
37 () => { "{\"capability_id\":\"lenso.configuration.source@1\",\"descriptor_version\":\"1.0.0\",\"cardinality\":\"many\"}" };
38 ($requirement_id:literal) => { concat!("{\"requirement_id\":", stringify!($requirement_id), ",\"capability_id\":\"lenso.configuration.source@1\",\"descriptor_version\":\"1.0.0\",\"cardinality\":\"many\"}") };
39}
40
41pub const FETCH_OPERATION: &str = "fetch";
42
43pub use lenso_contract_runtime::{Uint64, UnknownDomainError};
44use lenso_contract_runtime::{decode_portable_json, encode_portable_json};
45
46#[derive(Clone, Copy, Debug, PartialEq, serde::Serialize, serde::Deserialize)]
47pub struct FetchRequest {
48
49}
50
51#[derive(Clone, Debug, PartialEq, serde::Serialize, serde::Deserialize)]
52pub struct FetchResponse {
53 #[serde(rename = "configurations")]
54 #[serde(deserialize_with = "lenso_contract_runtime::serde::deserialize_required")]
55 pub configurations: Vec<FetchResponseConfigurationsItem>,
56 #[serde(rename = "revision")]
57 #[serde(deserialize_with = "lenso_contract_runtime::serde::deserialize_required")]
58 pub revision: Uint64,
59}
60
61#[derive(Clone, Debug, PartialEq, serde::Serialize, serde::Deserialize)]
62pub struct FetchResponseConfigurationsItem {
63 #[serde(rename = "instance_key")]
64 #[serde(deserialize_with = "lenso_contract_runtime::serde::deserialize_required")]
65 pub instance_key: String,
66 #[serde(rename = "plugin_id")]
67 #[serde(deserialize_with = "lenso_contract_runtime::serde::deserialize_required")]
68 pub plugin_id: String,
69 #[serde(rename = "toml")]
70 #[serde(deserialize_with = "lenso_contract_runtime::serde::deserialize_required")]
71 pub toml: String,
72}
73
74#[derive(Clone, Debug, PartialEq)]
75pub enum FetchError {
76 InvalidSource,
77 Unknown(UnknownDomainError),
78}
79
80#[derive(Debug)]
81pub struct Source;
82impl RequestCapability for Source {
83 type Request = FetchRequest;
84 type Response = FetchResponse;
85 type DomainError = FetchError;
86 const ID: &'static str = CAPABILITY_ID;
87 const DESCRIPTOR_VERSION: &'static str = DESCRIPTOR_VERSION;
88
89 fn invoke_native(endpoint: &dyn NativeRequestEndpoint, operation: &str, request: Self::Request, context: InvocationContext) -> NativeRequestFuture<Self> {
90 if operation != FETCH_OPERATION {
91 return lenso_kernel::invoke_typed_or_erased_native_request::<Self>(endpoint, operation, request, context);
92 }
93 let Some(typed_endpoint) = endpoint
94 .typed_endpoint()
95 .and_then(|endpoint| endpoint.downcast_ref::<SourceRequestEndpoint>())
96 else {
97 return lenso_kernel::invoke_typed_or_erased_native_request::<Self>(endpoint, operation, request, context);
98 };
99 Rc::clone(&typed_endpoint.provider).fetch(context, request)
100 }
101}
102
103impl serde::Serialize for FetchError {
104 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
105 where
106 S: serde::Serializer,
107 {
108 use serde::ser::SerializeMap;
109 match self {
110 Self::InvalidSource => serializer.serialize_str("invalid_source"),
111 Self::Unknown(value) => {
112 let mut map = serializer.serialize_map(Some(1 + usize::from(value.payload.is_some()) + value.extra.len()))?;
113 map.serialize_entry("code", &value.code)?;
114 if let Some(payload) = &value.payload {
115 map.serialize_entry("payload", payload)?;
116 }
117 for (key, extra) in &value.extra {
118 map.serialize_entry(key, extra)?;
119 }
120 map.end()
121 },
122 }
123 }
124}
125
126impl<'de> serde::Deserialize<'de> for FetchError {
127 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
128 where
129 D: serde::Deserializer<'de>,
130 {
131 let value = <serde_json::Value as serde::Deserialize>::deserialize(deserializer)?;
132 match value {
133 serde_json::Value::String(code) => match code.as_str() {
134 "invalid_source" => Ok(Self::InvalidSource),
135 _ => Ok(Self::Unknown(UnknownDomainError { code, payload: None, extra: std::collections::BTreeMap::new() })),
136 },
137 serde_json::Value::Object(mut object) => {
138 let Some(code) = object.remove("code").and_then(|value| value.as_str().map(ToOwned::to_owned)) else {
139 return Err(serde::de::Error::custom("Domain Error object is missing a string code"));
140 };
141 let payload = object.remove("payload");
142 let extra = object.into_iter().collect::<std::collections::BTreeMap<_, _>>();
143 Ok(Self::Unknown(UnknownDomainError { code, payload, extra }))
144 }
145 other => Err(serde::de::Error::custom(format!("Domain Error must be a string or object, got {other}"))),
146 }
147 }
148}
149
150pub fn encode_fetch_request(value: &FetchRequest) -> Result<String, serde_json::Error> { encode_portable_json(value) }
151pub fn decode_fetch_request(wire: &str) -> Result<FetchRequest, serde_json::Error> { decode_portable_json(wire) }
152pub fn encode_fetch_response(value: &FetchResponse) -> Result<String, serde_json::Error> { encode_portable_json(value) }
153pub fn decode_fetch_response(wire: &str) -> Result<FetchResponse, serde_json::Error> { decode_portable_json(wire) }
154pub fn encode_fetch_error(value: &FetchError) -> Result<String, serde_json::Error> { encode_portable_json(value) }
155pub fn decode_fetch_error(wire: &str) -> Result<FetchError, serde_json::Error> { decode_portable_json(wire) }
156
157#[doc(hidden)]
158pub trait __LensoIntoSourceFetchResult {
159 fn __lenso_into_result(self) -> Result<Result<FetchResponse, FetchError>, RuntimeFailure>;
160}
161impl __LensoIntoSourceFetchResult for Result<FetchResponse, FetchError> {
162 fn __lenso_into_result(self) -> Result<Result<FetchResponse, FetchError>, RuntimeFailure> { Ok(self) }
163}
164impl __LensoIntoSourceFetchResult for Result<Result<FetchResponse, FetchError>, RuntimeFailure> {
165 fn __lenso_into_result(self) -> Result<Result<FetchResponse, FetchError>, RuntimeFailure> { self }
166}
167impl __LensoIntoSourceFetchResult for Result<FetchResponse, lenso_plugin_authoring::PluginError<FetchError, RuntimeFailure>> {
168 fn __lenso_into_result(self) -> Result<Result<FetchResponse, FetchError>, RuntimeFailure> {
169 match self {
170 Ok(value) => Ok(Ok(value)),
171 Err(lenso_plugin_authoring::PluginError::Domain(error)) => Ok(Err(error)),
172 Err(lenso_plugin_authoring::PluginError::Runtime(error)) => Err(error),
173 }
174 }
175}
176impl __LensoIntoSourceFetchResult for Result<FetchResponse, SourceInvocationError> {
177 fn __lenso_into_result(self) -> Result<Result<FetchResponse, FetchError>, RuntimeFailure> {
178 match self {
179 Ok(value) => Ok(Ok(value)),
180 Err(SourceInvocationError::Domain(error)) => Ok(Err(error)),
181 Err(SourceInvocationError::Runtime(error)) => Err(error),
182 }
183 }
184}
185
186pub trait SourceProvider: fmt::Debug + 'static {
187 fn fetch(&self, context: InvocationContext, request: FetchRequest) -> NativeRequestFuture<Source>;
188}
189
190#[doc(hidden)]
191#[macro_export]
192macro_rules! __lenso_native_lower_source {
193 ($plugin:ty, $support:path) => {
194 use $support as __LensoNativeSupportSource;
195 impl $crate::SourceProvider for $plugin {
196 fn fetch(&self, context: __LensoNativeSupportSource::InvocationContext, request: $crate::FetchRequest) -> __LensoNativeSupportSource::NativeRequestFuture<$crate::Source> {
197 let plugin = self.clone();
198 ::std::boxed::Box::pin(async move {
199 let result = <$plugin>::fetch(&plugin, context, request).await;
200 $crate::__LensoIntoSourceFetchResult::__lenso_into_result(result)
201 })
202 }
203 }
204 };
205}
206
207#[doc(hidden)]
208#[macro_export]
209macro_rules! __lenso_native_lower_object_source {
210 ($object:ty, $plugin:ty, $support:path) => {
211 use $support as __LensoNativeSupportSource;
212 impl $crate::SourceProvider for $object {
213 fn fetch(&self, context: __LensoNativeSupportSource::InvocationContext, request: $crate::FetchRequest) -> __LensoNativeSupportSource::NativeRequestFuture<$crate::Source> {
214 let object = self.clone();
215 ::std::boxed::Box::pin(async move {
216 let plugin = object.get()?;
217 let result = <$plugin>::fetch(plugin.as_ref(), context, request).await;
218 $crate::__LensoIntoSourceFetchResult::__lenso_into_result(result)
219 })
220 }
221 }
222 };
223}
224
225#[doc(hidden)]
226#[macro_export]
227macro_rules! __lenso_native_lower_trait_object_source {
228 ($object:ty, $plugin:ty, $support:path) => {
229 use $support as __LensoNativeSupportSource;
230 impl $crate::SourceProvider for $object {
231 fn fetch(&self, context: __LensoNativeSupportSource::InvocationContext, request: $crate::FetchRequest) -> __LensoNativeSupportSource::NativeRequestFuture<$crate::Source> {
232 let object = self.clone();
233 ::std::boxed::Box::pin(async move {
234 let plugin = object.get()?;
235 <$plugin as $crate::SourceProvider>::fetch(plugin.as_ref(), context, request).await
236 })
237 }
238 }
239 };
240}
241
242#[derive(Debug)]
243struct SourceRequestEndpoint { provider: Rc<dyn SourceProvider> }
244
245#[derive(Debug)]
246pub struct SourceEndpoint<P: SourceProvider> { provider: Rc<P>, request_endpoint: SourceRequestEndpoint }
247impl<P: SourceProvider> SourceEndpoint<P> {
248 pub fn new(provider: P) -> Self {
249 let provider = Rc::new(provider);
250 let request_provider: Rc<dyn SourceProvider> = provider.clone();
251 Self { provider, request_endpoint: SourceRequestEndpoint { provider: request_provider } }
252 }
253}
254
255impl<P: SourceProvider> NativeRequestEndpoint for SourceEndpoint<P> {
256 fn capability_id(&self) -> &'static str { CAPABILITY_ID }
257 fn descriptor_version(&self) -> &'static str { DESCRIPTOR_VERSION }
258 fn operations(&self) -> &'static [&'static str] { &[
259 FETCH_OPERATION,
260 ] }
261 fn typed_endpoint(&self) -> Option<&dyn std::any::Any> { Some(&self.request_endpoint) }
262 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>> {
263 match operation {
264 FETCH_OPERATION => {
265 let Ok(request) = request.downcast::<FetchRequest>() else {
266 return Box::pin(futures::future::ready(Err(RuntimeFailure::ProtocolViolation { capability: CAPABILITY_ID })));
267 };
268 let invocation = Rc::clone(&self.provider).fetch(context, *request);
269 Box::pin(async move {
270 invocation.await.map(|result| {
271 result
272 .map(|value| Box::new(value) as Box<dyn std::any::Any>)
273 .map_err(|error| Box::new(error) as Box<dyn std::any::Any>)
274 })
275 })
276 }
277 _ => Box::pin(futures::future::ready(Err(RuntimeFailure::UnknownOperation { capability: CAPABILITY_ID, operation: operation.to_owned() }))),
278 }
279 }
280}
281
282#[doc(hidden)]
283#[macro_export]
284macro_rules! __lenso_native_endpoints_source {
285 ($provider:expr, $support:path) => {{
286 use $support as __LensoNativeSupport;
287 let endpoint = ::std::rc::Rc::new($crate::SourceEndpoint::new($provider));
288 (
289 vec![endpoint.clone() as ::std::rc::Rc<dyn __LensoNativeSupport::NativeRequestEndpoint>],
290 vec![],
291 vec![],
292 )
293 }};
294}
295
296#[doc(hidden)]
297#[macro_export]
298macro_rules! __lenso_native_provide_source {
299 ($provider:expr, $lifecycle:expr, $support:path) => {{
300 use $support as __LensoNativeSupport;
301 let (request_endpoints, stream_endpoints, event_endpoints) =
302 $crate::__lenso_native_endpoints_source!($provider, $support);
303 __LensoNativeSupport::NativePluginInstance::with_all_endpoints(
304 request_endpoints,
305 stream_endpoints,
306 event_endpoints,
307 $lifecycle,
308 )
309 }};
310}
311
312#[derive(Clone, Debug)]
313pub struct SourceClient {
314 fetch: NativeRequestHandle<Source>,
315}
316impl SourceClient {
317 pub fn new(handle: NativeRequestHandle<Source>) -> Self {
318 Self { fetch: handle }
319 }
320
321 pub fn from_dependencies(dependencies: &PluginDependencies) -> Result<Self, RuntimeFailure> {
322 <Self as CapabilityClient>::from_dependencies(dependencies)
323 }
324
325 pub fn from_requirement(
326 dependencies: &PluginDependencies,
327 requirement_id: &str,
328 ) -> Result<Self, RuntimeFailure> {
329 <Self as CapabilityClient>::from_requirement(dependencies, requirement_id)
330 }
331
332 pub async fn fetch(&self, request: FetchRequest) -> Result<FetchResponse, SourceInvocationError> {
333 self.fetch.invoke(FETCH_OPERATION, request).await
334 .map_err(SourceInvocationError::Runtime)?
335 .map_err(SourceInvocationError::Domain)
336 }
337
338 pub async fn fetch_with_context(&self, context: InvocationContext, request: FetchRequest) -> Result<FetchResponse, SourceInvocationError> {
339 self.fetch.invoke_with_context(FETCH_OPERATION, context, request).await
340 .map_err(SourceInvocationError::Runtime)?
341 .map_err(SourceInvocationError::Domain)
342 }
343}
344
345impl CapabilityClient for SourceClient {
346 type Dependencies = PluginDependencies;
347 type Error = RuntimeFailure;
348
349 const CAPABILITY_ID: &'static str = CAPABILITY_ID;
350 const DESCRIPTOR_VERSION: &'static str = DESCRIPTOR_VERSION;
351
352 fn from_dependencies(dependencies: &PluginDependencies) -> Result<Self, RuntimeFailure> {
353 Ok(Self {
354 fetch: dependencies.one::<Source>()?,
355 })
356 }
357
358 fn from_requirement(
359 dependencies: &PluginDependencies,
360 requirement_id: &str,
361 ) -> Result<Self, RuntimeFailure> {
362 let dependencies = dependencies.requirement(requirement_id)?;
363 Self::from_dependencies(&dependencies)
364 }
365
366 fn already_connected() -> RuntimeFailure {
367 RuntimeFailure::PluginFailure {
368 detail: format!("Capability Port {CAPABILITY_ID} was connected more than once"),
369 }
370 }
371}
372
373impl CapabilityClientMany for SourceClient {
374 fn many_from_dependencies(
375 dependencies: &PluginDependencies,
376 ) -> Result<Vec<BoundCapabilityClient<Self>>, RuntimeFailure> {
377 dependencies
378 .bindings()
379 .iter()
380 .filter(|binding| binding.capability_id() == CAPABILITY_ID)
381 .map(|binding| {
382 Ok(BoundCapabilityClient::new(
383 binding.provider_instance(),
384 Self {
385 fetch: binding.handle().ok_or(RuntimeFailure::Unavailable { capability: CAPABILITY_ID })?.typed::<Source>()?,
386 },
387 ))
388 })
389 .collect()
390 }
391
392 fn many_from_requirement(
393 dependencies: &PluginDependencies,
394 requirement_id: &str,
395 ) -> Result<Vec<BoundCapabilityClient<Self>>, RuntimeFailure> {
396 let dependencies = dependencies.requirement(requirement_id)?;
397 Self::many_from_dependencies(&dependencies)
398 }
399}
400
401#[derive(Clone, Debug, PartialEq)]
402pub enum SourceInvocationError {
403 Domain(FetchError),
404 Runtime(RuntimeFailure),
405}
406
407#[derive(Clone, Copy, Debug)]
408pub struct SourceGuestClient<'a, H: lenso_guest_sdk::HostImports> {
409 capability: lenso_guest_sdk::GuestCapability<'a, H>,
410}
411
412impl<'a, H: lenso_guest_sdk::HostImports> SourceGuestClient<'a, H> {
413 pub fn from_context(context: &'a lenso_guest_sdk::GuestContext<H>) -> Result<Self, lenso_guest_sdk::GuestError<serde_json::Value>> {
414 context
415 .require(CAPABILITY_ID, DESCRIPTOR_VERSION, &[FETCH_OPERATION], &[], &[])
416 .map(|capability| Self { capability })
417 }
418
419 pub fn fetch(&self, request: &FetchRequest) -> Result<FetchResponse, lenso_guest_sdk::GuestError<FetchError>> {
420 self.capability.request(FETCH_OPERATION, request)
421 }
422}
423
424#[derive(Debug, Default)]
425pub struct SourceJsonCodec;
426
427impl lenso_runtime_codec::JsonCapabilityCodec for SourceJsonCodec {
428 fn capability_id(&self) -> &'static str { CAPABILITY_ID }
429
430 fn descriptor_version(&self) -> &'static str { DESCRIPTOR_VERSION }
431
432 fn descriptor_digest(&self) -> &'static str { DESCRIPTOR_DIGEST }
433
434 fn request_operations(&self) -> &'static [&'static str] { &[FETCH_OPERATION] }
435 fn stream_operations(&self) -> &'static [&'static str] { &[] }
436 fn event_operations(&self) -> &'static [&'static str] { &[] }
437
438 fn encode_request(&self, operation: &str, request: &dyn std::any::Any) -> Result<serde_json::Value, RuntimeFailure> {
439 match operation {
440 FETCH_OPERATION => {
441 let value = request.downcast_ref::<FetchRequest>().ok_or_else(runtime_codec_protocol_failure)?;
442 serde_json::to_value(value).map_err(|_| runtime_codec_protocol_failure())
443 },
444 _ => Err(runtime_codec_unknown_operation(operation)),
445 }
446 }
447
448 fn decode_response(&self, operation: &str, value: serde_json::Value) -> Result<Box<dyn std::any::Any>, RuntimeFailure> {
449 match operation {
450 FETCH_OPERATION => serde_json::from_value::<FetchResponse>(value)
451 .map(|value| Box::new(value) as Box<dyn std::any::Any>)
452 .map_err(|_| runtime_codec_protocol_failure()),
453 _ => Err(runtime_codec_unknown_operation(operation)),
454 }
455 }
456
457 fn decode_domain_error(&self, operation: &str, value: serde_json::Value) -> Result<Box<dyn std::any::Any>, RuntimeFailure> {
458 match operation {
459 FETCH_OPERATION => serde_json::from_value::<FetchError>(value)
460 .map(|value| Box::new(value) as Box<dyn std::any::Any>)
461 .map_err(|_| runtime_codec_protocol_failure()),
462 _ => Err(runtime_codec_unknown_operation(operation)),
463 }
464 }
465
466 fn encode_stream_open(&self, operation: &str, _request: &dyn std::any::Any) -> Result<serde_json::Value, RuntimeFailure> {
467 Err(runtime_codec_unknown_operation(operation))
468 }
469
470 fn encode_stream_message(&self, operation: &str, _message: &dyn std::any::Any) -> Result<serde_json::Value, RuntimeFailure> {
471 Err(runtime_codec_unknown_operation(operation))
472 }
473
474 fn decode_stream_message(&self, operation: &str, _value: serde_json::Value) -> Result<Box<dyn std::any::Any>, RuntimeFailure> {
475 Err(runtime_codec_unknown_operation(operation))
476 }
477
478 fn decode_stream_domain_error(&self, operation: &str, _value: serde_json::Value) -> Result<Box<dyn std::any::Any>, RuntimeFailure> {
479 Err(runtime_codec_unknown_operation(operation))
480 }
481
482 fn encode_event(&self, operation: &str, _event: &dyn std::any::Any) -> Result<serde_json::Value, RuntimeFailure> {
483 Err(runtime_codec_unknown_operation(operation))
484 }
485
486 fn invoke_host_request(&self, dependency: lenso_kernel::PluginDependencyHandle, operation: String, request: serde_json::Value, context: InvocationContext) -> lenso_runtime_codec::JsonHostRequestFuture {
487 match operation.as_str() {
488 FETCH_OPERATION => {
489 let request = serde_json::from_value::<FetchRequest>(request).map_err(|_| runtime_codec_protocol_failure());
490 Box::pin(async move {
491 let request = request?;
492 let handle = dependency.typed::<Source>()?;
493 match handle.invoke_with_context(FETCH_OPERATION, context, request).await? {
494 Ok(response) => serde_json::to_value(response)
495 .map(lenso_runtime_codec::JsonInvocationOutcome::Success)
496 .map_err(|_| runtime_codec_protocol_failure()),
497 Err(error) => serde_json::to_value(error)
498 .map(lenso_runtime_codec::JsonInvocationOutcome::DomainError)
499 .map_err(|_| runtime_codec_protocol_failure()),
500 }
501 })
502 },
503 _ => Box::pin(std::future::ready(Err(runtime_codec_unknown_operation(&operation)))),
504 }
505 }
506
507 fn open_host_stream(&self, _dependency: lenso_kernel::PluginStreamDependencyHandle, operation: String, _request: serde_json::Value, _context: InvocationContext) -> lenso_runtime_codec::JsonHostStreamOpenFuture {
508 Box::pin(std::future::ready(Err(runtime_codec_unknown_operation(&operation))))
509 }
510
511 fn publish_host_event(&self, _dependency: lenso_kernel::PluginEventDependencyHandle, operation: String, _event: serde_json::Value, _context: InvocationContext) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>> {
512 Box::pin(std::future::ready(Err(runtime_codec_unknown_operation(&operation))))
513 }
514}
515
516fn runtime_codec_protocol_failure() -> RuntimeFailure { RuntimeFailure::ProtocolViolation { capability: CAPABILITY_ID } }
517
518fn runtime_codec_unknown_operation(operation: &str) -> RuntimeFailure {
519 RuntimeFailure::UnknownOperation { capability: CAPABILITY_ID, operation: operation.to_owned() }
520}