1use std::{
4 any::Any,
5 collections::BTreeMap,
6 env, fs,
7 io::Write as _,
8 path::{Path, PathBuf},
9 rc::Rc,
10 sync::Arc,
11};
12
13use lenso_app_plan::{
14 CapabilityCardinality, ExecutionClassId, PluginInstancePlan, ResolvedAppPlan,
15};
16use lenso_kernel::{
17 InvocationContext, NativeEventEndpoint, NativeRequestEndpoint, NativeStream,
18 NativeStreamEndpoint, NativeStreamItem, NativeStreamSession, PluginDependencies,
19 PluginDependencyHandle, PluginEventDependencyHandle, PluginStreamDependencyHandle,
20 PreparedBinding, PreparedNativeApp, PreparedNativePlugin, PreparedStreamBinding,
21 RuntimeFailure, StreamCapability, StreamEvent,
22};
23use serde::{Deserialize, Serialize};
24use serde_json::Value;
25use sha2::{Digest, Sha256};
26
27#[derive(Clone, Debug, Eq, PartialEq)]
29pub struct InstanceResources {
30 digest: String,
31 total_size: u64,
32 files: BTreeMap<String, Arc<[u8]>>,
33}
34
35impl Default for InstanceResources {
36 fn default() -> Self {
37 Self::from_files([]).expect("the empty resource snapshot is valid")
38 }
39}
40
41impl InstanceResources {
42 pub fn from_files(
44 files: impl IntoIterator<Item = (String, Vec<u8>)>,
45 ) -> Result<Self, RuntimeFailure> {
46 let mut indexed = BTreeMap::<String, Arc<[u8]>>::new();
47 let mut total_size = 0_u64;
48 for (path, bytes) in files {
49 validate_resource_path(&path)?;
50 total_size = total_size
51 .checked_add(
52 u64::try_from(bytes.len())
53 .map_err(|_| invalid_resources("resource file is too large"))?,
54 )
55 .ok_or_else(|| invalid_resources("resource snapshot size overflow"))?;
56 if indexed.insert(path.clone(), Arc::from(bytes)).is_some() {
57 return Err(invalid_resources(format!(
58 "duplicate Plugin resource path `{path}`"
59 )));
60 }
61 }
62 let mut hasher = Sha256::new();
63 hasher.update(b"lenso.instance-resources@1\0");
64 for (path, bytes) in &indexed {
65 hasher.update(
66 u64::try_from(path.len())
67 .expect("path length fits u64")
68 .to_be_bytes(),
69 );
70 hasher.update(path.as_bytes());
71 hasher.update(
72 u64::try_from(bytes.len())
73 .expect("content length fits u64")
74 .to_be_bytes(),
75 );
76 hasher.update(bytes.as_ref());
77 }
78 Ok(Self {
79 digest: format!("sha256:{}", hex::encode(hasher.finalize())),
80 total_size,
81 files: indexed,
82 })
83 }
84
85 pub fn digest(&self) -> &str {
87 &self.digest
88 }
89
90 pub fn file_count(&self) -> usize {
92 self.files.len()
93 }
94
95 pub const fn total_size(&self) -> u64 {
97 self.total_size
98 }
99
100 pub fn paths(&self) -> impl Iterator<Item = &str> {
102 self.files.keys().map(String::as_str)
103 }
104
105 pub fn read(&self, path: &str) -> Result<&[u8], RuntimeFailure> {
107 validate_resource_path(path)?;
108 self.files
109 .get(path)
110 .map(AsRef::as_ref)
111 .ok_or_else(|| invalid_resources(format!("Plugin resource `{path}` was not found")))
112 }
113
114 pub fn read_text(&self, path: &str) -> Result<&str, RuntimeFailure> {
116 std::str::from_utf8(self.read(path)?)
117 .map_err(|_| invalid_resources(format!("Plugin resource `{path}` is not UTF-8")))
118 }
119}
120
121#[derive(Clone, Debug, Default)]
123pub struct InstanceResourceCatalog {
124 snapshots: BTreeMap<String, InstanceResources>,
125 empty: InstanceResources,
126}
127
128impl InstanceResourceCatalog {
129 pub fn new() -> Self {
131 Self::default()
132 }
133
134 pub fn with_resources(
136 mut self,
137 instance_key: impl Into<String>,
138 resources: InstanceResources,
139 ) -> Result<Self, RuntimeFailure> {
140 let instance_key = instance_key.into();
141 if self
142 .snapshots
143 .insert(instance_key.clone(), resources)
144 .is_some()
145 {
146 return Err(invalid_resources(format!(
147 "duplicate resource authority for Instance `{instance_key}`"
148 )));
149 }
150 Ok(self)
151 }
152
153 pub fn for_instance(&self, instance_key: &str) -> &InstanceResources {
155 self.snapshots.get(instance_key).unwrap_or(&self.empty)
156 }
157
158 pub fn iter(&self) -> impl Iterator<Item = (&str, &InstanceResources)> {
160 self.snapshots
161 .iter()
162 .map(|(instance, resources)| (instance.as_str(), resources))
163 }
164}
165
166#[derive(Debug)]
168struct ArtifactBacking {
169 path: PathBuf,
170 _directory: tempfile::TempDir,
171}
172
173#[derive(Clone)]
175pub struct ArtifactHandle {
176 source_path: PathBuf,
177 backing: Arc<ArtifactBacking>,
178 digest: String,
179 size: u64,
180}
181
182impl std::fmt::Debug for ArtifactHandle {
183 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
184 formatter
185 .debug_struct("ArtifactHandle")
186 .field("source_path", &self.source_path)
187 .field("path", &self.backing.path)
188 .field("digest", &self.digest)
189 .field("size", &self.size)
190 .finish_non_exhaustive()
191 }
192}
193
194impl PartialEq for ArtifactHandle {
195 fn eq(&self, other: &Self) -> bool {
196 self.source_path == other.source_path
197 && self.digest == other.digest
198 && self.size == other.size
199 }
200}
201
202impl Eq for ArtifactHandle {}
203
204impl ArtifactHandle {
205 pub fn open(
210 path: impl Into<PathBuf>,
211 expected_digest: &str,
212 expected_size: u64,
213 ) -> Result<Self, RuntimeFailure> {
214 Self::open_inner(path.into(), expected_digest, expected_size, None)
215 }
216
217 pub fn open_with_staging_root(
220 path: impl Into<PathBuf>,
221 expected_digest: &str,
222 expected_size: u64,
223 staging_root: impl AsRef<Path>,
224 ) -> Result<Self, RuntimeFailure> {
225 Self::open_inner(
226 path.into(),
227 expected_digest,
228 expected_size,
229 Some(staging_root.as_ref()),
230 )
231 }
232
233 fn open_inner(
234 path: PathBuf,
235 expected_digest: &str,
236 expected_size: u64,
237 staging_root: Option<&Path>,
238 ) -> Result<Self, RuntimeFailure> {
239 validate_digest(expected_digest)?;
240 let path = absolute_path(path)?;
241 let metadata =
242 fs::symlink_metadata(&path).map_err(|error| invalid_artifact(&path, error))?;
243 if !metadata.file_type().is_file() || metadata.file_type().is_symlink() {
244 return Err(RuntimeFailure::InvalidResolvedPlan {
245 detail: format!("Artifact `{}` is not a regular file", path.display()),
246 });
247 }
248 if metadata.len() != expected_size {
249 return Err(RuntimeFailure::InvalidResolvedPlan {
250 detail: format!(
251 "Artifact `{}` size mismatch: expected {expected_size}, got {}",
252 path.display(),
253 metadata.len()
254 ),
255 });
256 }
257 let mut source = fs::File::open(&path).map_err(|error| invalid_artifact(&path, error))?;
258 let opened_metadata = source
259 .metadata()
260 .map_err(|error| invalid_artifact(&path, error))?;
261 if !opened_metadata.is_file() || opened_metadata.len() != expected_size {
262 return Err(RuntimeFailure::InvalidResolvedPlan {
263 detail: format!("Artifact `{}` changed during admission", path.display()),
264 });
265 }
266 let (backing, actual_digest, actual_size) =
267 materialize_stable_artifact(&path, &mut source, &opened_metadata, staging_root)?;
268 if actual_size != expected_size {
269 return Err(RuntimeFailure::InvalidResolvedPlan {
270 detail: format!("Artifact `{}` changed during admission", path.display()),
271 });
272 }
273 if actual_digest != expected_digest {
274 return Err(RuntimeFailure::InvalidResolvedPlan {
275 detail: format!("Artifact `{}` digest mismatch", path.display()),
276 });
277 }
278 Ok(Self {
279 source_path: path,
280 backing: Arc::new(backing),
281 digest: actual_digest,
282 size: opened_metadata.len(),
283 })
284 }
285
286 pub fn path(&self) -> &Path {
289 &self.backing.path
290 }
291
292 pub fn source_path(&self) -> &Path {
294 &self.source_path
295 }
296
297 pub fn digest(&self) -> &str {
299 &self.digest
300 }
301
302 pub const fn size(&self) -> u64 {
304 self.size
305 }
306
307 pub fn read_verified(&self) -> Result<Vec<u8>, RuntimeFailure> {
309 let bytes = fs::read(&self.backing.path)
310 .map_err(|error| invalid_artifact(&self.backing.path, error))?;
311 let size = u64::try_from(bytes.len()).unwrap_or(u64::MAX);
312 let digest = format!("sha256:{}", hex::encode(Sha256::digest(&bytes)));
313 if size != self.size || digest != self.digest {
314 return Err(RuntimeFailure::InvalidResolvedPlan {
315 detail: format!(
316 "stable Artifact `{}` changed after admission",
317 self.backing.path.display()
318 ),
319 });
320 }
321 Ok(bytes)
322 }
323}
324
325fn absolute_path(path: PathBuf) -> Result<PathBuf, RuntimeFailure> {
326 if path.is_absolute() {
327 return Ok(path);
328 }
329 env::current_dir()
330 .map(|current| current.join(path))
331 .map_err(|error| invalid_artifact(Path::new("."), error))
332}
333
334fn materialize_stable_artifact(
335 source_path: &Path,
336 source: &mut fs::File,
337 source_metadata: &fs::Metadata,
338 staging_root: Option<&Path>,
339) -> Result<(ArtifactBacking, String, u64), RuntimeFailure> {
340 let directory = stable_artifact_directory(source_path, staging_root)?;
341 let file_name = source_path
342 .file_name()
343 .ok_or_else(|| RuntimeFailure::InvalidResolvedPlan {
344 detail: format!("Artifact `{}` has no filename", source_path.display()),
345 })?;
346 let stable_path = directory.path().join(file_name);
347 let mut stable = fs::OpenOptions::new()
348 .create_new(true)
349 .write(true)
350 .open(&stable_path)
351 .map_err(|error| invalid_artifact(source_path, error))?;
352 let mut hasher = Sha256::new();
353 let mut size = 0_u64;
354 let mut buffer = vec![0_u8; 64 * 1024];
355 loop {
356 let read = std::io::Read::read(source, &mut buffer)
357 .map_err(|error| invalid_artifact(source_path, error))?;
358 if read == 0 {
359 break;
360 }
361 size = size
362 .checked_add(u64::try_from(read).expect("buffer length fits u64"))
363 .ok_or_else(|| RuntimeFailure::InvalidResolvedPlan {
364 detail: format!("Artifact `{}` size overflow", source_path.display()),
365 })?;
366 hasher.update(&buffer[..read]);
367 stable
368 .write_all(&buffer[..read])
369 .map_err(|error| invalid_artifact(source_path, error))?;
370 }
371 set_stable_permissions(&stable_path, source_metadata)
372 .map_err(|error| invalid_artifact(source_path, error))?;
373 Ok((
374 ArtifactBacking {
375 path: stable_path,
376 _directory: directory,
377 },
378 format!("sha256:{}", hex::encode(hasher.finalize())),
379 size,
380 ))
381}
382
383fn stable_artifact_directory(
384 source_path: &Path,
385 staging_root: Option<&Path>,
386) -> Result<tempfile::TempDir, RuntimeFailure> {
387 let builder = || {
388 let mut builder = tempfile::Builder::new();
389 builder.prefix("lenso-artifact-");
390 builder
391 };
392 if let Some(root) = staging_root {
393 return builder()
394 .tempdir_in(root)
395 .map_err(|error| invalid_artifact(source_path, error));
396 }
397 builder()
398 .tempdir()
399 .map_err(|error| invalid_artifact(source_path, error))
400}
401
402#[cfg(unix)]
403fn set_stable_permissions(path: &Path, metadata: &fs::Metadata) -> std::io::Result<()> {
404 use std::os::unix::fs::PermissionsExt as _;
405
406 let mode = metadata.permissions().mode() & 0o555;
407 fs::set_permissions(path, fs::Permissions::from_mode(mode))
408}
409
410#[cfg(not(unix))]
411fn set_stable_permissions(path: &Path, metadata: &fs::Metadata) -> std::io::Result<()> {
412 let mut permissions = metadata.permissions();
413 permissions.set_readonly(true);
414 fs::set_permissions(path, permissions)
415}
416
417#[derive(Clone, Debug, Default)]
419pub struct ArtifactCatalog(BTreeMap<String, ArtifactHandle>);
420
421impl ArtifactCatalog {
422 pub fn new() -> Self {
424 Self::default()
425 }
426
427 pub fn with_artifact(
429 mut self,
430 instance_key: impl Into<String>,
431 artifact: ArtifactHandle,
432 ) -> Result<Self, RuntimeFailure> {
433 let instance_key = instance_key.into();
434 if self.0.insert(instance_key.clone(), artifact).is_some() {
435 return Err(RuntimeFailure::InvalidResolvedPlan {
436 detail: format!("duplicate Artifact authority for Instance `{instance_key}`"),
437 });
438 }
439 Ok(self)
440 }
441
442 pub fn require(&self, instance_key: &str) -> Result<&ArtifactHandle, RuntimeFailure> {
444 self.0
445 .get(instance_key)
446 .ok_or_else(|| RuntimeFailure::InvalidResolvedPlan {
447 detail: format!("no admitted Artifact for Instance `{instance_key}`"),
448 })
449 }
450}
451
452pub trait JsonCapabilityCodec: std::fmt::Debug + 'static {
454 fn capability_id(&self) -> &'static str;
456 fn descriptor_version(&self) -> &'static str;
458 fn descriptor_digest(&self) -> &'static str {
464 ""
465 }
466 fn request_operations(&self) -> &'static [&'static str];
468 fn stream_operations(&self) -> &'static [&'static str] {
470 &[]
471 }
472 fn event_operations(&self) -> &'static [&'static str] {
474 &[]
475 }
476 fn encode_request(&self, operation: &str, request: &dyn Any) -> Result<Value, RuntimeFailure>;
478 fn decode_response(
480 &self,
481 operation: &str,
482 value: Value,
483 ) -> Result<Box<dyn Any>, RuntimeFailure>;
484 fn decode_domain_error(
486 &self,
487 operation: &str,
488 value: Value,
489 ) -> Result<Box<dyn Any>, RuntimeFailure>;
490 fn encode_stream_open(
492 &self,
493 operation: &str,
494 request: &dyn Any,
495 ) -> Result<Value, RuntimeFailure> {
496 let _ = request;
497 Err(unknown_operation(self.capability_id(), operation))
498 }
499 fn encode_stream_message(
501 &self,
502 operation: &str,
503 message: &dyn Any,
504 ) -> Result<Value, RuntimeFailure> {
505 let _ = message;
506 Err(unknown_operation(self.capability_id(), operation))
507 }
508 fn decode_stream_message(
510 &self,
511 operation: &str,
512 value: Value,
513 ) -> Result<Box<dyn Any>, RuntimeFailure> {
514 let _ = value;
515 Err(unknown_operation(self.capability_id(), operation))
516 }
517 fn decode_stream_domain_error(
519 &self,
520 operation: &str,
521 value: Value,
522 ) -> Result<Box<dyn Any>, RuntimeFailure> {
523 let _ = value;
524 Err(unknown_operation(self.capability_id(), operation))
525 }
526 fn encode_event(&self, operation: &str, event: &dyn Any) -> Result<Value, RuntimeFailure> {
528 let _ = event;
529 Err(unknown_operation(self.capability_id(), operation))
530 }
531 fn invoke_host_request(
533 &self,
534 dependency: PluginDependencyHandle,
535 operation: String,
536 request: Value,
537 context: InvocationContext,
538 ) -> JsonHostRequestFuture {
539 let _ = (dependency, request, context);
540 Box::pin(futures::future::ready(Err(unknown_operation(
541 self.capability_id(),
542 &operation,
543 ))))
544 }
545 fn open_host_stream(
547 &self,
548 dependency: PluginStreamDependencyHandle,
549 operation: String,
550 request: Value,
551 context: InvocationContext,
552 ) -> JsonHostStreamOpenFuture {
553 let _ = (dependency, request, context);
554 Box::pin(futures::future::ready(Err(unknown_operation(
555 self.capability_id(),
556 &operation,
557 ))))
558 }
559 fn publish_host_event(
561 &self,
562 dependency: PluginEventDependencyHandle,
563 operation: String,
564 event: Value,
565 context: InvocationContext,
566 ) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>> {
567 let _ = (dependency, event, context);
568 Box::pin(futures::future::ready(Err(unknown_operation(
569 self.capability_id(),
570 &operation,
571 ))))
572 }
573}
574
575#[derive(Debug)]
577pub enum JsonInvocationOutcome {
578 Success(Value),
580 DomainError(Value),
582}
583
584pub fn json_runtime_failure(error: &RuntimeFailure) -> Value {
586 match error {
587 RuntimeFailure::Unavailable { capability } => serde_json::json!({
588 "kind": "unavailable",
589 "capability": capability,
590 }),
591 RuntimeFailure::UnknownOperation {
592 capability,
593 operation,
594 } => serde_json::json!({
595 "kind": "unknown_operation",
596 "capability": capability,
597 "operation": operation,
598 }),
599 RuntimeFailure::AmbiguousBinding {
600 capability,
601 providers,
602 } => serde_json::json!({
603 "kind": "ambiguous_binding",
604 "capability": capability,
605 "providers": providers,
606 }),
607 RuntimeFailure::ProtocolViolation { capability } => serde_json::json!({
608 "kind": "protocol_violation",
609 "capability": capability,
610 }),
611 RuntimeFailure::AdmissionClosed => serde_json::json!({ "kind": "admission_closed" }),
612 RuntimeFailure::ResourceExhausted {
613 capability,
614 operation,
615 } => serde_json::json!({
616 "kind": "resource_exhausted",
617 "capability": capability,
618 "operation": operation,
619 }),
620 RuntimeFailure::DeadlineExceeded { request_id } => serde_json::json!({
621 "kind": "deadline_exceeded",
622 "request_id": request_id.to_string(),
623 }),
624 RuntimeFailure::Cancelled { request_id } => serde_json::json!({
625 "kind": "cancelled",
626 "request_id": request_id.to_string(),
627 }),
628 RuntimeFailure::MissingPluginFactory { .. }
629 | RuntimeFailure::UnavailableExecutionClass { .. }
630 | RuntimeFailure::InvalidResolvedPlan { .. }
631 | RuntimeFailure::Internal { .. }
632 | RuntimeFailure::PluginFailure { .. }
633 | RuntimeFailure::PluginRestartExhausted { .. } => {
634 serde_json::json!({ "kind": "internal" })
635 }
636 }
637}
638
639pub fn json_host_invocation_envelope(
641 outcome: Result<JsonInvocationOutcome, RuntimeFailure>,
642) -> Value {
643 match outcome {
644 Ok(JsonInvocationOutcome::Success(value)) => serde_json::json!({ "ok": value }),
645 Ok(JsonInvocationOutcome::DomainError(value)) => serde_json::json!({ "error": value }),
646 Err(error) => serde_json::json!({ "runtime": json_runtime_failure(&error) }),
647 }
648}
649
650pub type JsonHostRequestFuture =
652 futures::future::LocalBoxFuture<'static, Result<JsonInvocationOutcome, RuntimeFailure>>;
653
654pub trait JsonHostStreamSession: std::fmt::Debug + 'static {
656 fn send(
657 self: Rc<Self>,
658 message: Value,
659 ) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>>;
660 fn receive(
661 self: Rc<Self>,
662 ) -> futures::future::LocalBoxFuture<'static, Result<JsonStreamItem, RuntimeFailure>>;
663 fn close_send(
664 self: Rc<Self>,
665 ) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>>;
666 fn cancel(&self);
667}
668
669pub type JsonHostStreamOpenFuture = futures::future::LocalBoxFuture<
671 'static,
672 Result<Result<Rc<dyn JsonHostStreamSession>, Value>, RuntimeFailure>,
673>;
674
675type DecodeStreamMessage<C> =
676 Rc<dyn Fn(Value) -> Result<<C as StreamCapability>::Message, RuntimeFailure>>;
677type EncodeStreamMessage<C> =
678 Rc<dyn Fn(<C as StreamCapability>::Message) -> Result<Value, RuntimeFailure>>;
679type EncodeStreamError<C> =
680 Rc<dyn Fn(<C as StreamCapability>::DomainError) -> Result<Value, RuntimeFailure>>;
681
682pub fn json_host_stream<C: StreamCapability>(
684 stream: NativeStream<C>,
685 decode_message: impl Fn(Value) -> Result<C::Message, RuntimeFailure> + 'static,
686 encode_message: impl Fn(C::Message) -> Result<Value, RuntimeFailure> + 'static,
687 encode_error: impl Fn(C::DomainError) -> Result<Value, RuntimeFailure> + 'static,
688) -> Rc<dyn JsonHostStreamSession> {
689 Rc::new(TypedJsonHostStream {
690 stream: Rc::new(stream),
691 decode_message: Rc::new(decode_message),
692 encode_message: Rc::new(encode_message),
693 encode_error: Rc::new(encode_error),
694 })
695}
696
697struct TypedJsonHostStream<C: StreamCapability> {
698 stream: Rc<NativeStream<C>>,
699 decode_message: DecodeStreamMessage<C>,
700 encode_message: EncodeStreamMessage<C>,
701 encode_error: EncodeStreamError<C>,
702}
703
704impl<C: StreamCapability> std::fmt::Debug for TypedJsonHostStream<C> {
705 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
706 formatter
707 .debug_struct("TypedJsonHostStream")
708 .field("capability", &C::ID)
709 .finish_non_exhaustive()
710 }
711}
712
713impl<C: StreamCapability> JsonHostStreamSession for TypedJsonHostStream<C> {
714 fn send(
715 self: Rc<Self>,
716 message: Value,
717 ) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>> {
718 Box::pin(async move {
719 let message = (self.decode_message)(message)?;
720 self.stream.send(message).await
721 })
722 }
723
724 fn receive(
725 self: Rc<Self>,
726 ) -> futures::future::LocalBoxFuture<'static, Result<JsonStreamItem, RuntimeFailure>> {
727 Box::pin(async move {
728 match self.stream.receive().await? {
729 StreamEvent::Message(message) => {
730 (self.encode_message)(message).map(JsonStreamItem::Message)
731 }
732 StreamEvent::PeerHalfClosed => Ok(JsonStreamItem::PeerHalfClosed),
733 StreamEvent::Terminal(Ok(())) => Ok(JsonStreamItem::Terminal(Ok(()))),
734 StreamEvent::Terminal(Err(error)) => {
735 (self.encode_error)(error).map(|error| JsonStreamItem::Terminal(Err(error)))
736 }
737 }
738 })
739 }
740
741 fn close_send(
742 self: Rc<Self>,
743 ) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>> {
744 Box::pin(async move { self.stream.close_send().await })
745 }
746
747 fn cancel(&self) {
748 self.stream.cancel();
749 }
750}
751
752pub const JSON_REQUEST_ABI_V1: &str = "lenso.json-request@1";
754
755pub const JSON_INTERACTIONS_ABI_V1: &str = "lenso.json-interactions@1";
757
758pub const JSON_HOST_IMPORTS_ABI_V1: &str = "lenso.json-host-imports@1";
760pub const JSON_HOST_IMPORTS_ABI_V2: &str = "lenso.json-host-imports@2";
762
763#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
765#[serde(deny_unknown_fields)]
766pub struct JsonPluginDescriptor {
767 pub abi: String,
768 pub capabilities: Vec<JsonCapabilityDescriptor>,
769 #[serde(default, skip_serializing_if = "Vec::is_empty")]
770 pub required_capabilities: Vec<JsonRequiredCapabilityDescriptor>,
771}
772
773#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
775#[serde(deny_unknown_fields)]
776pub struct JsonCapabilityDescriptor {
777 pub capability_id: String,
778 pub descriptor_version: String,
779 pub request_operations: Vec<String>,
780 #[serde(default, skip_serializing_if = "Vec::is_empty")]
781 pub stream_operations: Vec<String>,
782}
783
784#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
786#[serde(deny_unknown_fields)]
787pub struct JsonRequiredCapabilityDescriptor {
788 pub requirement_id: String,
789 pub capability_id: String,
790 pub descriptor_version: String,
791 pub cardinality: CapabilityCardinality,
792}
793
794pub fn expected_json_plugin_descriptor(
796 instance: &PluginInstancePlan,
797) -> Result<JsonPluginDescriptor, RuntimeFailure> {
798 let mut capabilities = Vec::with_capacity(instance.provided_capabilities().len());
799 for descriptor in instance.provided_capabilities() {
800 if !descriptor.event_operations().is_empty() {
801 return Err(RuntimeFailure::InvalidResolvedPlan {
802 detail: format!(
803 "Execution class `{}` does not support Event endpoints through the legacy JSON descriptor",
804 instance.execution_class()
805 ),
806 });
807 }
808 capabilities.push(JsonCapabilityDescriptor {
809 capability_id: descriptor.capability_id().to_owned(),
810 descriptor_version: descriptor.descriptor_version().to_owned(),
811 request_operations: descriptor
812 .request_operations()
813 .into_iter()
814 .map(str::to_owned)
815 .collect(),
816 stream_operations: descriptor
817 .stream_operations()
818 .into_iter()
819 .map(str::to_owned)
820 .collect(),
821 });
822 }
823 capabilities.sort();
824 if capabilities
825 .windows(2)
826 .any(|pair| pair[0].capability_id == pair[1].capability_id)
827 {
828 return Err(RuntimeFailure::InvalidResolvedPlan {
829 detail: format!(
830 "Instance `{}` declares a duplicate Capability",
831 instance.instance_key()
832 ),
833 });
834 }
835 let mut required_capabilities = instance
836 .required_capabilities()
837 .iter()
838 .map(|requirement| JsonRequiredCapabilityDescriptor {
839 requirement_id: requirement.requirement_id().to_owned(),
840 capability_id: requirement.capability_id().to_owned(),
841 descriptor_version: requirement.descriptor_version().to_owned(),
842 cardinality: requirement.cardinality(),
843 })
844 .collect::<Vec<_>>();
845 sort_required_capabilities(&mut required_capabilities);
846 Ok(JsonPluginDescriptor {
847 abi: if !required_capabilities.is_empty() {
848 JSON_HOST_IMPORTS_ABI_V2
849 } else if capabilities
850 .iter()
851 .any(|capability| !capability.stream_operations.is_empty())
852 {
853 JSON_INTERACTIONS_ABI_V1
854 } else {
855 JSON_REQUEST_ABI_V1
856 }
857 .to_owned(),
858 capabilities,
859 required_capabilities,
860 })
861}
862
863pub fn validate_json_plugin_descriptor(
865 instance: &PluginInstancePlan,
866 encoded: &str,
867) -> Result<(), RuntimeFailure> {
868 let mut actual = serde_json::from_str::<JsonPluginDescriptor>(encoded).map_err(|_| {
869 RuntimeFailure::ProtocolViolation {
870 capability: "lenso.json-request@1",
871 }
872 })?;
873 actual.capabilities.sort();
874 sort_required_capabilities(&mut actual.required_capabilities);
875 let expected = expected_json_plugin_descriptor(instance)?;
876 if actual != expected {
877 return Err(RuntimeFailure::InvalidResolvedPlan {
878 detail: format!(
879 "guest descriptor does not match resolved Instance `{}`",
880 instance.instance_key()
881 ),
882 });
883 }
884 Ok(())
885}
886
887fn sort_required_capabilities(requirements: &mut [JsonRequiredCapabilityDescriptor]) {
888 requirements.sort_by(|left, right| {
889 (
890 &left.requirement_id,
891 &left.capability_id,
892 &left.descriptor_version,
893 cardinality_order(left.cardinality),
894 )
895 .cmp(&(
896 &right.requirement_id,
897 &right.capability_id,
898 &right.descriptor_version,
899 cardinality_order(right.cardinality),
900 ))
901 });
902}
903
904const fn cardinality_order(cardinality: CapabilityCardinality) -> u8 {
905 match cardinality {
906 CapabilityCardinality::One => 0,
907 CapabilityCardinality::Optional => 1,
908 CapabilityCardinality::Many => 2,
909 }
910}
911
912pub trait JsonRequestTransport: std::fmt::Debug + 'static {
914 fn invoke(
915 self: Rc<Self>,
916 capability: String,
917 operation: String,
918 request_json: String,
919 context: InvocationContext,
920 ) -> futures::future::LocalBoxFuture<'static, Result<JsonInvocationOutcome, RuntimeFailure>>;
921}
922
923#[derive(Debug)]
925pub enum JsonStreamItem {
926 Message(Value),
927 PeerHalfClosed,
928 Terminal(Result<(), Value>),
929}
930
931#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
933#[serde(
934 tag = "kind",
935 content = "value",
936 rename_all = "kebab-case",
937 deny_unknown_fields
938)]
939pub enum JsonStreamFrame {
940 Message(Value),
941 PeerHalfClosed,
942 TerminalSuccess,
943 TerminalError(Value),
944}
945
946#[derive(Clone, Debug, Eq, PartialEq, Serialize)]
948pub struct JsonHostBindingDescriptor {
949 pub binding_id: u32,
950 pub requirement_id: String,
951 pub provider_instance: String,
952 pub capability_id: String,
953 pub descriptor_version: String,
954 #[serde(skip_serializing)]
955 pub descriptor_digest: String,
956 pub request_operations: Vec<String>,
957 pub stream_operations: Vec<String>,
958 #[serde(skip_serializing_if = "Vec::is_empty")]
959 pub event_operations: Vec<String>,
960}
961
962#[derive(Clone)]
963struct JsonHostBinding {
964 descriptor: JsonHostBindingDescriptor,
965 codec: Rc<dyn JsonCapabilityCodec>,
966 request: Option<PluginDependencyHandle>,
967 stream: Option<PluginStreamDependencyHandle>,
968 event: Option<PluginEventDependencyHandle>,
969}
970
971impl std::fmt::Debug for JsonHostBinding {
972 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
973 formatter
974 .debug_struct("JsonHostBinding")
975 .field("descriptor", &self.descriptor)
976 .finish_non_exhaustive()
977 }
978}
979
980#[derive(Debug)]
982pub struct JsonHostImports {
983 codecs: BTreeMap<String, Rc<dyn JsonCapabilityCodec>>,
984 bindings: std::cell::RefCell<Option<Vec<JsonHostBinding>>>,
985 streams: std::cell::RefCell<BTreeMap<u64, Rc<dyn JsonHostStreamSession>>>,
986 next_stream_id: std::cell::Cell<u64>,
987 max_streams: usize,
988}
989
990impl JsonHostImports {
991 pub fn new(
993 codecs: Vec<Rc<dyn JsonCapabilityCodec>>,
994 max_streams: usize,
995 ) -> Result<Self, RuntimeFailure> {
996 let mut by_capability = BTreeMap::new();
997 for codec in codecs {
998 let capability = codec.capability_id().to_owned();
999 if let Some(existing) = by_capability.get(&capability) {
1000 if !Rc::ptr_eq(existing, &codec) {
1001 return Err(RuntimeFailure::InvalidResolvedPlan {
1002 detail: format!(
1003 "conflicting guest import codecs for Capability `{capability}`"
1004 ),
1005 });
1006 }
1007 } else {
1008 by_capability.insert(capability, codec);
1009 }
1010 }
1011 Ok(Self {
1012 codecs: by_capability,
1013 bindings: std::cell::RefCell::new(None),
1014 streams: std::cell::RefCell::new(BTreeMap::new()),
1015 next_stream_id: std::cell::Cell::new(1),
1016 max_streams,
1017 })
1018 }
1019
1020 pub fn activate(&self, dependencies: &PluginDependencies) -> Result<(), RuntimeFailure> {
1022 if self.bindings.borrow().is_some() {
1023 return Err(RuntimeFailure::Internal {
1024 detail: "guest Capability imports were activated twice".to_owned(),
1025 });
1026 }
1027 let mut bindings = Vec::with_capacity(dependencies.len());
1028 for (index, dependency) in dependencies.bindings().iter().enumerate() {
1029 let codec = self
1030 .codecs
1031 .get(dependency.capability_id())
1032 .cloned()
1033 .ok_or_else(|| RuntimeFailure::InvalidResolvedPlan {
1034 detail: format!(
1035 "no generated guest import codec for Capability `{}`",
1036 dependency.capability_id()
1037 ),
1038 })?;
1039 let request = dependency.handle();
1040 let stream = dependency.stream_handle();
1041 let event = dependency.event_handle();
1042 validate_host_binding(&codec, request.as_ref(), stream.as_ref(), event.as_ref())?;
1043 let binding_id =
1044 u32::try_from(index).map_err(|_| RuntimeFailure::InvalidResolvedPlan {
1045 detail: "guest import binding table exceeds u32 identity space".to_owned(),
1046 })?;
1047 bindings.push(JsonHostBinding {
1048 descriptor: JsonHostBindingDescriptor {
1049 binding_id,
1050 requirement_id: dependency.requirement_id().to_owned(),
1051 provider_instance: dependency.provider_instance().to_owned(),
1052 capability_id: dependency.capability_id().to_owned(),
1053 descriptor_version: codec.descriptor_version().to_owned(),
1054 descriptor_digest: codec.descriptor_digest().to_owned(),
1055 request_operations: request.as_ref().map_or_else(Vec::new, |handle| {
1056 handle
1057 .operations()
1058 .iter()
1059 .map(|item| (*item).to_owned())
1060 .collect()
1061 }),
1062 stream_operations: stream.as_ref().map_or_else(Vec::new, |handle| {
1063 handle
1064 .operations()
1065 .iter()
1066 .map(|item| (*item).to_owned())
1067 .collect()
1068 }),
1069 event_operations: event.as_ref().map_or_else(Vec::new, |handle| {
1070 handle
1071 .operations()
1072 .iter()
1073 .map(|item| (*item).to_owned())
1074 .collect()
1075 }),
1076 },
1077 codec,
1078 request,
1079 stream,
1080 event,
1081 });
1082 }
1083 self.bindings.replace(Some(bindings));
1084 Ok(())
1085 }
1086
1087 pub fn descriptors(&self) -> Result<Vec<JsonHostBindingDescriptor>, RuntimeFailure> {
1089 self.bindings
1090 .borrow()
1091 .as_ref()
1092 .map(|bindings| {
1093 bindings
1094 .iter()
1095 .map(|binding| binding.descriptor.clone())
1096 .collect()
1097 })
1098 .ok_or(RuntimeFailure::AdmissionClosed)
1099 }
1100
1101 pub fn invoke(
1103 &self,
1104 binding_id: u32,
1105 operation: String,
1106 request: Value,
1107 context: InvocationContext,
1108 ) -> JsonHostRequestFuture {
1109 let binding = match self.binding(binding_id) {
1110 Ok(binding) => binding,
1111 Err(error) => return Box::pin(futures::future::ready(Err(error))),
1112 };
1113 let Some(dependency) = binding.request else {
1114 return Box::pin(futures::future::ready(Err(
1115 RuntimeFailure::UnknownOperation {
1116 capability: binding.codec.capability_id(),
1117 operation,
1118 },
1119 )));
1120 };
1121 binding
1122 .codec
1123 .invoke_host_request(dependency, operation, request, context)
1124 }
1125
1126 pub fn publish_event(
1128 &self,
1129 binding_id: u32,
1130 operation: String,
1131 event: Value,
1132 context: InvocationContext,
1133 ) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>> {
1134 let binding = match self.binding(binding_id) {
1135 Ok(binding) => binding,
1136 Err(error) => return Box::pin(futures::future::ready(Err(error))),
1137 };
1138 let Some(dependency) = binding.event else {
1139 return Box::pin(futures::future::ready(Err(
1140 RuntimeFailure::UnknownOperation {
1141 capability: binding.codec.capability_id(),
1142 operation,
1143 },
1144 )));
1145 };
1146 binding
1147 .codec
1148 .publish_host_event(dependency, operation, event, context)
1149 }
1150
1151 pub fn open_stream(
1153 self: Rc<Self>,
1154 binding_id: u32,
1155 operation: String,
1156 request: Value,
1157 context: InvocationContext,
1158 ) -> futures::future::LocalBoxFuture<'static, Result<Result<u64, Value>, RuntimeFailure>> {
1159 Box::pin(async move {
1160 if self.streams.borrow().len() >= self.max_streams {
1161 return Err(RuntimeFailure::ResourceExhausted {
1162 capability: JSON_HOST_IMPORTS_ABI_V2,
1163 operation: "stream-open".to_owned(),
1164 });
1165 }
1166 let binding = self.binding(binding_id)?;
1167 let dependency = binding
1168 .stream
1169 .ok_or_else(|| RuntimeFailure::UnknownOperation {
1170 capability: binding.codec.capability_id(),
1171 operation: operation.clone(),
1172 })?;
1173 match binding
1174 .codec
1175 .open_host_stream(dependency, operation, request, context)
1176 .await?
1177 {
1178 Ok(stream) => {
1179 let stream_id = self.next_stream_id.get();
1180 let next =
1181 stream_id
1182 .checked_add(1)
1183 .ok_or(RuntimeFailure::ResourceExhausted {
1184 capability: JSON_HOST_IMPORTS_ABI_V2,
1185 operation: "stream-open".to_owned(),
1186 })?;
1187 self.next_stream_id.set(next);
1188 self.streams.borrow_mut().insert(stream_id, stream);
1189 Ok(Ok(stream_id))
1190 }
1191 Err(error) => Ok(Err(error)),
1192 }
1193 })
1194 }
1195
1196 pub fn send_stream(
1198 &self,
1199 stream_id: u64,
1200 message: Value,
1201 ) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>> {
1202 match self.stream(stream_id) {
1203 Ok(stream) => stream.send(message),
1204 Err(error) => Box::pin(futures::future::ready(Err(error))),
1205 }
1206 }
1207
1208 pub fn receive_stream(
1210 self: Rc<Self>,
1211 stream_id: u64,
1212 ) -> futures::future::LocalBoxFuture<'static, Result<JsonStreamItem, RuntimeFailure>> {
1213 Box::pin(async move {
1214 let stream = self.stream(stream_id)?;
1215 let item = stream.receive().await?;
1216 if matches!(item, JsonStreamItem::Terminal(_)) {
1217 self.streams.borrow_mut().remove(&stream_id);
1218 }
1219 Ok(item)
1220 })
1221 }
1222
1223 pub fn close_stream_send(
1225 &self,
1226 stream_id: u64,
1227 ) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>> {
1228 match self.stream(stream_id) {
1229 Ok(stream) => stream.close_send(),
1230 Err(error) => Box::pin(futures::future::ready(Err(error))),
1231 }
1232 }
1233
1234 pub fn cancel_stream(&self, stream_id: u64) -> Result<(), RuntimeFailure> {
1236 let stream = self
1237 .streams
1238 .borrow_mut()
1239 .remove(&stream_id)
1240 .ok_or_else(unknown_host_stream)?;
1241 stream.cancel();
1242 Ok(())
1243 }
1244
1245 pub fn deactivate(&self) {
1247 self.bindings.replace(None);
1248 for (_, stream) in std::mem::take(&mut *self.streams.borrow_mut()) {
1249 stream.cancel();
1250 }
1251 }
1252
1253 fn binding(&self, binding_id: u32) -> Result<JsonHostBinding, RuntimeFailure> {
1254 let bindings = self.bindings.borrow();
1255 let bindings = bindings.as_ref().ok_or(RuntimeFailure::AdmissionClosed)?;
1256 bindings
1257 .get(binding_id as usize)
1258 .cloned()
1259 .ok_or(RuntimeFailure::ProtocolViolation {
1260 capability: JSON_HOST_IMPORTS_ABI_V2,
1261 })
1262 }
1263
1264 fn stream(&self, stream_id: u64) -> Result<Rc<dyn JsonHostStreamSession>, RuntimeFailure> {
1265 self.streams
1266 .borrow()
1267 .get(&stream_id)
1268 .cloned()
1269 .ok_or_else(unknown_host_stream)
1270 }
1271}
1272
1273fn validate_host_binding(
1274 codec: &Rc<dyn JsonCapabilityCodec>,
1275 request: Option<&PluginDependencyHandle>,
1276 stream: Option<&PluginStreamDependencyHandle>,
1277 event: Option<&PluginEventDependencyHandle>,
1278) -> Result<(), RuntimeFailure> {
1279 for (capability, version) in request
1280 .map(|handle| (handle.capability_id(), handle.descriptor_version()))
1281 .into_iter()
1282 .chain(stream.map(|handle| (handle.capability_id(), handle.descriptor_version())))
1283 .chain(event.map(|handle| (handle.capability_id(), handle.descriptor_version())))
1284 {
1285 if capability != codec.capability_id() || version != codec.descriptor_version() {
1286 return Err(RuntimeFailure::ProtocolViolation {
1287 capability: codec.capability_id(),
1288 });
1289 }
1290 }
1291 Ok(())
1292}
1293
1294fn unknown_host_stream() -> RuntimeFailure {
1295 RuntimeFailure::ProtocolViolation {
1296 capability: JSON_HOST_IMPORTS_ABI_V2,
1297 }
1298}
1299
1300impl JsonStreamFrame {
1301 pub fn decode(
1303 encoded: &str,
1304 capability: &'static str,
1305 ) -> Result<JsonStreamItem, RuntimeFailure> {
1306 match serde_json::from_str(encoded)
1307 .map_err(|_| RuntimeFailure::ProtocolViolation { capability })?
1308 {
1309 Self::Message(value) => Ok(JsonStreamItem::Message(value)),
1310 Self::PeerHalfClosed => Ok(JsonStreamItem::PeerHalfClosed),
1311 Self::TerminalSuccess => Ok(JsonStreamItem::Terminal(Ok(()))),
1312 Self::TerminalError(value) => Ok(JsonStreamItem::Terminal(Err(value))),
1313 }
1314 }
1315}
1316
1317pub trait JsonStreamSessionTransport: std::fmt::Debug + 'static {
1319 fn send(
1320 self: Rc<Self>,
1321 message_json: String,
1322 ) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>>;
1323 fn receive(
1324 self: Rc<Self>,
1325 ) -> futures::future::LocalBoxFuture<'static, Result<JsonStreamItem, RuntimeFailure>>;
1326 fn close_send(
1327 self: Rc<Self>,
1328 ) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>>;
1329 fn cancel(&self);
1330}
1331
1332pub type JsonStreamOpenFuture = futures::future::LocalBoxFuture<
1334 'static,
1335 Result<Result<Rc<dyn JsonStreamSessionTransport>, Value>, RuntimeFailure>,
1336>;
1337
1338pub trait JsonStreamTransport: std::fmt::Debug + 'static {
1340 fn open(
1341 self: Rc<Self>,
1342 capability: String,
1343 operation: String,
1344 request_json: String,
1345 context: InvocationContext,
1346 ) -> JsonStreamOpenFuture;
1347}
1348
1349pub trait JsonEventTransport: std::fmt::Debug + 'static {
1351 fn publish(
1353 self: Rc<Self>,
1354 capability: String,
1355 operation: String,
1356 event_json: String,
1357 context: InvocationContext,
1358 ) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>>;
1359}
1360
1361pub fn json_request_endpoints<T: JsonRequestTransport>(
1363 transport: Rc<T>,
1364 codecs: Vec<Rc<dyn JsonCapabilityCodec>>,
1365) -> Vec<Rc<dyn NativeRequestEndpoint>> {
1366 let transport: Rc<dyn JsonRequestTransport> = transport;
1367 codecs
1368 .into_iter()
1369 .filter(|codec| !codec.request_operations().is_empty())
1370 .map(|codec| {
1371 Rc::new(JsonRequestEndpoint {
1372 transport: transport.clone(),
1373 codec,
1374 }) as Rc<dyn NativeRequestEndpoint>
1375 })
1376 .collect()
1377}
1378
1379pub fn json_stream_endpoints<T: JsonStreamTransport>(
1381 transport: Rc<T>,
1382 codecs: Vec<Rc<dyn JsonCapabilityCodec>>,
1383) -> Vec<Rc<dyn NativeStreamEndpoint>> {
1384 let transport: Rc<dyn JsonStreamTransport> = transport;
1385 codecs
1386 .into_iter()
1387 .filter(|codec| !codec.stream_operations().is_empty())
1388 .map(|codec| {
1389 Rc::new(JsonStreamEndpoint {
1390 transport: transport.clone(),
1391 codec,
1392 }) as Rc<dyn NativeStreamEndpoint>
1393 })
1394 .collect()
1395}
1396
1397pub fn json_event_endpoints<T: JsonEventTransport>(
1399 transport: Rc<T>,
1400 codecs: Vec<Rc<dyn JsonCapabilityCodec>>,
1401) -> Vec<Rc<dyn NativeEventEndpoint>> {
1402 let transport: Rc<dyn JsonEventTransport> = transport;
1403 codecs
1404 .into_iter()
1405 .filter(|codec| !codec.event_operations().is_empty())
1406 .map(|codec| {
1407 Rc::new(JsonEventEndpoint {
1408 transport: transport.clone(),
1409 codec,
1410 }) as Rc<dyn NativeEventEndpoint>
1411 })
1412 .collect()
1413}
1414
1415#[derive(Debug)]
1416struct JsonEventEndpoint {
1417 transport: Rc<dyn JsonEventTransport>,
1418 codec: Rc<dyn JsonCapabilityCodec>,
1419}
1420
1421impl NativeEventEndpoint for JsonEventEndpoint {
1422 fn capability_id(&self) -> &'static str {
1423 self.codec.capability_id()
1424 }
1425
1426 fn descriptor_version(&self) -> &'static str {
1427 self.codec.descriptor_version()
1428 }
1429
1430 fn operations(&self) -> &'static [&'static str] {
1431 self.codec.event_operations()
1432 }
1433
1434 fn owns_event_admission(&self) -> bool {
1435 true
1436 }
1437
1438 fn publish(
1439 &self,
1440 operation: &str,
1441 event: Box<dyn Any>,
1442 context: InvocationContext,
1443 ) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>> {
1444 if !self.codec.event_operations().contains(&operation) {
1445 return Box::pin(futures::future::ready(Err(unknown_operation(
1446 self.codec.capability_id(),
1447 operation,
1448 ))));
1449 }
1450 let event = self.codec.encode_event(operation, event.as_ref());
1451 let capability = self.codec.capability_id();
1452 let operation = operation.to_owned();
1453 let transport = self.transport.clone();
1454 Box::pin(async move {
1455 let event_json = serde_json::to_string(&event?)
1456 .map_err(|_| RuntimeFailure::ProtocolViolation { capability })?;
1457 transport
1458 .publish(capability.to_owned(), operation, event_json, context)
1459 .await
1460 })
1461 }
1462}
1463
1464#[derive(Debug)]
1465struct JsonStreamEndpoint {
1466 transport: Rc<dyn JsonStreamTransport>,
1467 codec: Rc<dyn JsonCapabilityCodec>,
1468}
1469
1470impl NativeStreamEndpoint for JsonStreamEndpoint {
1471 fn capability_id(&self) -> &'static str {
1472 self.codec.capability_id()
1473 }
1474 fn descriptor_version(&self) -> &'static str {
1475 self.codec.descriptor_version()
1476 }
1477 fn operations(&self) -> &'static [&'static str] {
1478 self.codec.stream_operations()
1479 }
1480
1481 fn open(
1482 &self,
1483 operation: &str,
1484 request: Box<dyn Any>,
1485 context: InvocationContext,
1486 ) -> futures::future::LocalBoxFuture<
1487 'static,
1488 Result<Result<Box<dyn NativeStreamSession>, Box<dyn Any>>, RuntimeFailure>,
1489 > {
1490 let transport = self.transport.clone();
1491 let codec = self.codec.clone();
1492 let operation = operation.to_owned();
1493 Box::pin(async move {
1494 if !codec.stream_operations().contains(&operation.as_str()) {
1495 return Err(unknown_operation(codec.capability_id(), &operation));
1496 }
1497 let request = codec.encode_stream_open(&operation, request.as_ref())?;
1498 let request_json =
1499 serde_json::to_string(&request).map_err(|_| RuntimeFailure::ProtocolViolation {
1500 capability: codec.capability_id(),
1501 })?;
1502 match transport
1503 .open(
1504 codec.capability_id().to_owned(),
1505 operation.clone(),
1506 request_json,
1507 context,
1508 )
1509 .await?
1510 {
1511 Ok(session) => Ok(Ok(Box::new(JsonStreamSession {
1512 session,
1513 codec,
1514 operation,
1515 }) as Box<dyn NativeStreamSession>)),
1516 Err(error) => codec.decode_stream_domain_error(&operation, error).map(Err),
1517 }
1518 })
1519 }
1520}
1521
1522#[derive(Debug)]
1523struct JsonStreamSession {
1524 session: Rc<dyn JsonStreamSessionTransport>,
1525 codec: Rc<dyn JsonCapabilityCodec>,
1526 operation: String,
1527}
1528
1529impl NativeStreamSession for JsonStreamSession {
1530 fn send(
1531 &self,
1532 message: Box<dyn Any>,
1533 ) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>> {
1534 let encoded = self
1535 .codec
1536 .encode_stream_message(&self.operation, message.as_ref())
1537 .and_then(|value| {
1538 serde_json::to_string(&value).map_err(|_| RuntimeFailure::ProtocolViolation {
1539 capability: self.codec.capability_id(),
1540 })
1541 });
1542 let session = self.session.clone();
1543 Box::pin(async move { session.send(encoded?).await })
1544 }
1545
1546 fn receive(
1547 &self,
1548 ) -> futures::future::LocalBoxFuture<'static, Result<NativeStreamItem, RuntimeFailure>> {
1549 let session = self.session.clone();
1550 let codec = self.codec.clone();
1551 let operation = self.operation.clone();
1552 Box::pin(async move {
1553 match session.receive().await? {
1554 JsonStreamItem::Message(value) => codec
1555 .decode_stream_message(&operation, value)
1556 .map(NativeStreamItem::Message),
1557 JsonStreamItem::PeerHalfClosed => Ok(NativeStreamItem::PeerHalfClosed),
1558 JsonStreamItem::Terminal(Ok(())) => Ok(NativeStreamItem::Terminal(Ok(()))),
1559 JsonStreamItem::Terminal(Err(value)) => codec
1560 .decode_stream_domain_error(&operation, value)
1561 .map(|error| NativeStreamItem::Terminal(Err(error))),
1562 }
1563 })
1564 }
1565
1566 fn close_send(&self) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>> {
1567 self.session.clone().close_send()
1568 }
1569
1570 fn cancel(&self) {
1571 self.session.cancel();
1572 }
1573}
1574
1575#[derive(Debug)]
1576struct JsonRequestEndpoint {
1577 transport: Rc<dyn JsonRequestTransport>,
1578 codec: Rc<dyn JsonCapabilityCodec>,
1579}
1580
1581impl NativeRequestEndpoint for JsonRequestEndpoint {
1582 fn capability_id(&self) -> &'static str {
1583 self.codec.capability_id()
1584 }
1585
1586 fn descriptor_version(&self) -> &'static str {
1587 self.codec.descriptor_version()
1588 }
1589
1590 fn operations(&self) -> &'static [&'static str] {
1591 self.codec.request_operations()
1592 }
1593
1594 fn invoke(
1595 &self,
1596 operation: &str,
1597 request: Box<dyn Any>,
1598 context: InvocationContext,
1599 ) -> futures::future::LocalBoxFuture<
1600 'static,
1601 Result<Result<Box<dyn Any>, Box<dyn Any>>, RuntimeFailure>,
1602 > {
1603 let transport = self.transport.clone();
1604 let codec = self.codec.clone();
1605 let operation = operation.to_owned();
1606 Box::pin(async move {
1607 if !codec.request_operations().contains(&operation.as_str()) {
1608 return Err(RuntimeFailure::UnknownOperation {
1609 capability: codec.capability_id(),
1610 operation,
1611 });
1612 }
1613 let request = codec.encode_request(&operation, request.as_ref())?;
1614 let request =
1615 serde_json::to_string(&request).map_err(|_| RuntimeFailure::ProtocolViolation {
1616 capability: codec.capability_id(),
1617 })?;
1618 match transport
1619 .invoke(
1620 codec.capability_id().to_owned(),
1621 operation.clone(),
1622 request,
1623 context,
1624 )
1625 .await?
1626 {
1627 JsonInvocationOutcome::Success(value) => {
1628 codec.decode_response(&operation, value).map(Ok)
1629 }
1630 JsonInvocationOutcome::DomainError(value) => {
1631 codec.decode_domain_error(&operation, value).map(Err)
1632 }
1633 }
1634 })
1635 }
1636}
1637
1638pub fn codecs_for_instance(
1640 instance: &PluginInstancePlan,
1641 codecs: &BTreeMap<String, Rc<dyn JsonCapabilityCodec>>,
1642) -> Result<Vec<Rc<dyn JsonCapabilityCodec>>, RuntimeFailure> {
1643 let mut selected = Vec::with_capacity(instance.provided_capabilities().len());
1644 for descriptor in instance.provided_capabilities() {
1645 let codec = codecs.get(descriptor.capability_id()).ok_or_else(|| {
1646 RuntimeFailure::InvalidResolvedPlan {
1647 detail: format!(
1648 "no generated codec for Capability `{}`",
1649 descriptor.capability_id()
1650 ),
1651 }
1652 })?;
1653 let request_operations: Vec<_> = codec
1654 .request_operations()
1655 .iter()
1656 .map(|operation| (*operation).to_owned())
1657 .collect();
1658 let stream_operations: Vec<_> = codec
1659 .stream_operations()
1660 .iter()
1661 .map(|operation| (*operation).to_owned())
1662 .collect();
1663 let event_operations: Vec<_> = codec
1664 .event_operations()
1665 .iter()
1666 .map(|operation| (*operation).to_owned())
1667 .collect();
1668 let expected_request: Vec<_> = descriptor
1669 .request_operations()
1670 .into_iter()
1671 .map(str::to_owned)
1672 .collect();
1673 let expected_stream: Vec<_> = descriptor
1674 .stream_operations()
1675 .into_iter()
1676 .map(str::to_owned)
1677 .collect();
1678 let expected_event: Vec<_> = descriptor
1679 .event_operations()
1680 .into_iter()
1681 .map(str::to_owned)
1682 .collect();
1683 if codec.descriptor_version() != descriptor.descriptor_version()
1684 || request_operations != expected_request
1685 || stream_operations != expected_stream
1686 || event_operations != expected_event
1687 {
1688 return Err(RuntimeFailure::ProtocolViolation {
1689 capability: codec.capability_id(),
1690 });
1691 }
1692 selected.push(codec.clone());
1693 }
1694 Ok(selected)
1695}
1696
1697pub fn codecs_for_requirements(
1699 instance: &PluginInstancePlan,
1700 codecs: &BTreeMap<String, Rc<dyn JsonCapabilityCodec>>,
1701) -> Result<Vec<Rc<dyn JsonCapabilityCodec>>, RuntimeFailure> {
1702 let mut selected = BTreeMap::new();
1703 for requirement in instance.required_capabilities() {
1704 let codec = codecs.get(requirement.capability_id()).ok_or_else(|| {
1705 RuntimeFailure::InvalidResolvedPlan {
1706 detail: format!(
1707 "no generated guest import codec for Capability `{}`",
1708 requirement.capability_id()
1709 ),
1710 }
1711 })?;
1712 if codec.descriptor_version() != requirement.descriptor_version() {
1713 return Err(RuntimeFailure::ProtocolViolation {
1714 capability: codec.capability_id(),
1715 });
1716 }
1717 selected
1718 .entry(requirement.capability_id().to_owned())
1719 .or_insert_with(|| codec.clone());
1720 }
1721 Ok(selected.into_values().collect())
1722}
1723
1724pub fn prepare_request_app(
1726 plan: &ResolvedAppPlan,
1727 execution_class: &ExecutionClassId,
1728 generations: BTreeMap<String, PreparedNativePlugin>,
1729) -> Result<PreparedNativeApp, RuntimeFailure> {
1730 let selected_instances = plan
1731 .plugin_instances()
1732 .iter()
1733 .filter(|instance| instance.execution_class() == execution_class)
1734 .map(|instance| instance.instance_key().to_owned())
1735 .collect::<std::collections::BTreeSet<_>>();
1736 let mut endpoints = BTreeMap::new();
1737 let mut stream_endpoints = BTreeMap::new();
1738 for (instance_key, generation) in &generations {
1739 for endpoint in generation.endpoints() {
1740 let identity = (instance_key.clone(), endpoint.capability_id().to_owned());
1741 if endpoints.insert(identity, endpoint.clone()).is_some() {
1742 return Err(RuntimeFailure::InvalidResolvedPlan {
1743 detail: format!("duplicate request endpoint on Instance `{instance_key}`"),
1744 });
1745 }
1746 }
1747 for endpoint in generation.stream_endpoints() {
1748 let identity = (instance_key.clone(), endpoint.capability_id().to_owned());
1749 if stream_endpoints
1750 .insert(identity, endpoint.clone())
1751 .is_some()
1752 {
1753 return Err(RuntimeFailure::InvalidResolvedPlan {
1754 detail: format!("duplicate stream endpoint on Instance `{instance_key}`"),
1755 });
1756 }
1757 }
1758 }
1759 for instance in plan
1760 .plugin_instances()
1761 .iter()
1762 .filter(|instance| selected_instances.contains(instance.instance_key()))
1763 {
1764 if !generations.contains_key(instance.instance_key()) {
1765 return Err(RuntimeFailure::InvalidResolvedPlan {
1766 detail: format!("Adapter omitted Instance `{}`", instance.instance_key()),
1767 });
1768 }
1769 }
1770 let mut bindings = Vec::new();
1771 let mut stream_bindings = Vec::new();
1772 for binding in plan.capability_bindings() {
1773 let key = (
1774 binding.provider_instance().to_owned(),
1775 binding.capability_id().to_owned(),
1776 );
1777 let request_endpoint = endpoints.get(&key);
1778 let stream_endpoint = stream_endpoints.get(&key);
1779 if let Some(endpoint) = request_endpoint {
1780 bindings.push(
1781 PreparedBinding::new(
1782 binding.consumer_instance(),
1783 binding.provider_instance(),
1784 endpoint.clone(),
1785 )
1786 .with_requirement_id(binding.requirement_id()),
1787 );
1788 }
1789 if let Some(endpoint) = stream_endpoint {
1790 stream_bindings.push(
1791 PreparedStreamBinding::new(
1792 binding.consumer_instance(),
1793 binding.provider_instance(),
1794 endpoint.clone(),
1795 )
1796 .with_requirement_id(binding.requirement_id()),
1797 );
1798 }
1799 if request_endpoint.is_none()
1800 && stream_endpoint.is_none()
1801 && selected_instances.contains(binding.provider_instance())
1802 {
1803 return Err(RuntimeFailure::InvalidResolvedPlan {
1804 detail: format!(
1805 "Adapter omitted Capability `{}` endpoint for Instance `{}`",
1806 binding.capability_id(),
1807 binding.provider_instance()
1808 ),
1809 });
1810 }
1811 }
1812 Ok(PreparedNativeApp::new(bindings, generations).with_stream_bindings(stream_bindings))
1813}
1814
1815pub fn require_operation(
1817 codecs: &BTreeMap<String, Rc<dyn JsonCapabilityCodec>>,
1818 capability_id: &str,
1819 operation: &str,
1820) -> Result<Rc<dyn JsonCapabilityCodec>, RuntimeFailure> {
1821 let codec =
1822 codecs
1823 .get(capability_id)
1824 .cloned()
1825 .ok_or_else(|| RuntimeFailure::InvalidResolvedPlan {
1826 detail: format!("no generated codec for Capability `{capability_id}`"),
1827 })?;
1828 if !codec.request_operations().contains(&operation) {
1829 return Err(RuntimeFailure::UnknownOperation {
1830 capability: codec.capability_id(),
1831 operation: operation.to_owned(),
1832 });
1833 }
1834 Ok(codec)
1835}
1836
1837fn unknown_operation(capability: &'static str, operation: &str) -> RuntimeFailure {
1838 RuntimeFailure::UnknownOperation {
1839 capability,
1840 operation: operation.to_owned(),
1841 }
1842}
1843
1844fn validate_digest(digest: &str) -> Result<(), RuntimeFailure> {
1845 let valid = digest.strip_prefix("sha256:").is_some_and(|hex| {
1846 hex.len() == 64
1847 && hex
1848 .bytes()
1849 .all(|byte| byte.is_ascii_hexdigit() && !byte.is_ascii_uppercase())
1850 });
1851 if valid {
1852 Ok(())
1853 } else {
1854 Err(RuntimeFailure::InvalidResolvedPlan {
1855 detail: format!("invalid canonical SHA-256 digest `{digest}`"),
1856 })
1857 }
1858}
1859
1860fn invalid_artifact(path: &Path, error: impl std::fmt::Display) -> RuntimeFailure {
1861 RuntimeFailure::InvalidResolvedPlan {
1862 detail: format!("cannot read Artifact `{}`: {error}", path.display()),
1863 }
1864}
1865
1866fn validate_resource_path(path: &str) -> Result<(), RuntimeFailure> {
1867 if path.is_empty()
1868 || path.starts_with('/')
1869 || path.contains(['\\', '\0'])
1870 || path
1871 .split('/')
1872 .any(|segment| segment.is_empty() || matches!(segment, "." | ".."))
1873 {
1874 return Err(invalid_resources(format!(
1875 "invalid Plugin resource path `{path}`"
1876 )));
1877 }
1878 Ok(())
1879}
1880
1881fn invalid_resources(detail: impl Into<String>) -> RuntimeFailure {
1882 RuntimeFailure::InvalidResolvedPlan {
1883 detail: detail.into(),
1884 }
1885}
1886
1887#[cfg(test)]
1888mod tests {
1889 use std::{cell::RefCell, io::Write};
1890
1891 use super::*;
1892
1893 #[derive(Debug)]
1894 struct EventCodec;
1895
1896 impl JsonCapabilityCodec for EventCodec {
1897 fn capability_id(&self) -> &'static str {
1898 "example.events@1"
1899 }
1900
1901 fn descriptor_version(&self) -> &'static str {
1902 "1.0.0"
1903 }
1904
1905 fn request_operations(&self) -> &'static [&'static str] {
1906 &[]
1907 }
1908
1909 fn event_operations(&self) -> &'static [&'static str] {
1910 &["changed"]
1911 }
1912
1913 fn encode_request(
1914 &self,
1915 operation: &str,
1916 _request: &dyn Any,
1917 ) -> Result<Value, RuntimeFailure> {
1918 Err(unknown_operation(self.capability_id(), operation))
1919 }
1920
1921 fn decode_response(
1922 &self,
1923 operation: &str,
1924 _value: Value,
1925 ) -> Result<Box<dyn Any>, RuntimeFailure> {
1926 Err(unknown_operation(self.capability_id(), operation))
1927 }
1928
1929 fn decode_domain_error(
1930 &self,
1931 operation: &str,
1932 _value: Value,
1933 ) -> Result<Box<dyn Any>, RuntimeFailure> {
1934 Err(unknown_operation(self.capability_id(), operation))
1935 }
1936
1937 fn encode_event(&self, operation: &str, event: &dyn Any) -> Result<Value, RuntimeFailure> {
1938 if operation != "changed" {
1939 return Err(unknown_operation(self.capability_id(), operation));
1940 }
1941 event
1942 .downcast_ref::<String>()
1943 .map(|value| serde_json::json!({"value": value}))
1944 .ok_or(RuntimeFailure::ProtocolViolation {
1945 capability: self.capability_id(),
1946 })
1947 }
1948 }
1949
1950 #[derive(Debug, Default)]
1951 struct EventTransport {
1952 published: RefCell<Option<(String, String, String)>>,
1953 }
1954
1955 impl JsonEventTransport for EventTransport {
1956 fn publish(
1957 self: Rc<Self>,
1958 capability: String,
1959 operation: String,
1960 event_json: String,
1961 _context: InvocationContext,
1962 ) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>> {
1963 Box::pin(async move {
1964 self.published
1965 .replace(Some((capability, operation, event_json)));
1966 Ok(())
1967 })
1968 }
1969 }
1970
1971 #[test]
1972 fn event_endpoints_encode_typed_values_and_commit_transport_admission() {
1973 let transport = Rc::new(EventTransport::default());
1974 let endpoints = json_event_endpoints(
1975 transport.clone(),
1976 vec![Rc::new(EventCodec) as Rc<dyn JsonCapabilityCodec>],
1977 );
1978 assert_eq!(endpoints.len(), 1);
1979 assert!(endpoints[0].owns_event_admission());
1980 assert_eq!(endpoints[0].operations(), &["changed"]);
1981
1982 futures::executor::block_on(endpoints[0].publish(
1983 "changed",
1984 Box::new("ready".to_owned()),
1985 InvocationContext::new(1, None, lenso_kernel::CancellationToken::new()),
1986 ))
1987 .unwrap();
1988
1989 assert_eq!(
1990 transport.published.borrow().as_ref().unwrap(),
1991 &(
1992 "example.events@1".to_owned(),
1993 "changed".to_owned(),
1994 r#"{"value":"ready"}"#.to_owned(),
1995 )
1996 );
1997 }
1998
1999 #[test]
2000 fn artifact_handle_keeps_the_admitted_bytes_after_source_drift() {
2001 let mut file = tempfile::NamedTempFile::new().unwrap();
2002 file.write_all(b"first").unwrap();
2003 let digest = format!("sha256:{}", hex::encode(Sha256::digest(b"first")));
2004 let handle = ArtifactHandle::open(file.path(), &digest, 5).unwrap();
2005 file.as_file_mut().set_len(0).unwrap();
2006 file.write_all(b"other").unwrap();
2007 assert_eq!(handle.read_verified().unwrap(), b"first");
2008 assert_ne!(handle.path(), file.path());
2009 assert_eq!(fs::read(handle.path()).unwrap(), b"first");
2010 }
2011
2012 #[test]
2013 fn artifact_snapshot_survives_source_parent_rename_and_replacement() {
2014 let workspace = tempfile::tempdir().unwrap();
2015 let selected = workspace.path().join("selected");
2016 fs::create_dir(&selected).unwrap();
2017 let source = selected.join("plugin");
2018 fs::write(&source, b"admitted").unwrap();
2019 let digest = format!("sha256:{}", hex::encode(Sha256::digest(b"admitted")));
2020
2021 let handle = ArtifactHandle::open(&source, &digest, 8).unwrap();
2022 assert!(!handle.path().starts_with(&selected));
2023 fs::rename(&selected, workspace.path().join("replaced")).unwrap();
2024 fs::create_dir(&selected).unwrap();
2025 fs::write(selected.join("plugin"), b"attacker").unwrap();
2026
2027 assert_eq!(fs::read(handle.path()).unwrap(), b"admitted");
2028 assert_eq!(handle.read_verified().unwrap(), b"admitted");
2029 }
2030
2031 #[test]
2032 fn artifact_admission_streams_large_content_into_one_stable_snapshot() {
2033 let bytes = vec![0x5a; 4 * 1024 * 1024 + 17];
2034 let mut file = tempfile::NamedTempFile::new().unwrap();
2035 file.write_all(&bytes).unwrap();
2036 let digest = format!("sha256:{}", hex::encode(Sha256::digest(&bytes)));
2037
2038 let handle = ArtifactHandle::open(file.path(), &digest, bytes.len() as u64).unwrap();
2039
2040 assert_eq!(handle.read_verified().unwrap(), bytes);
2041 }
2042
2043 #[test]
2046 #[ignore = "large Artifact admission benchmark; run explicitly"]
2047 fn artifact_admission_streaming_benchmark() {
2048 const BLOCK_BYTES: usize = 64 * 1024;
2049 let block = vec![0x5a; BLOCK_BYTES];
2050 for mebibytes in [4_usize, 64, 256] {
2051 let directory = tempfile::tempdir().unwrap();
2052 let path = directory.path().join("artifact");
2053 let mut source = fs::File::create(&path).unwrap();
2054 let mut hasher = Sha256::new();
2055 let blocks = mebibytes * 1024 * 1024 / BLOCK_BYTES;
2056 for _ in 0..blocks {
2057 source.write_all(&block).unwrap();
2058 hasher.update(&block);
2059 }
2060 drop(source);
2061 let size = u64::try_from(mebibytes * 1024 * 1024).unwrap();
2062 let digest = format!("sha256:{}", hex::encode(hasher.finalize()));
2063
2064 let started = std::time::Instant::now();
2065 let handle = ArtifactHandle::open(&path, &digest, size).unwrap();
2066 let elapsed = started.elapsed();
2067
2068 assert_eq!(handle.size(), size);
2069 println!(
2070 "{{\"mebibytes\":{mebibytes},\"elapsed_ms\":{:.3},\"mib_per_second\":{:.3}}}",
2071 elapsed.as_secs_f64() * 1_000.0,
2072 f64::from(u32::try_from(mebibytes).unwrap()) / elapsed.as_secs_f64()
2073 );
2074 drop(handle);
2075 }
2076 }
2077
2078 #[test]
2079 fn artifact_admission_honors_an_explicit_host_staging_root() {
2080 let source = tempfile::tempdir().unwrap();
2081 let staging = tempfile::tempdir().unwrap();
2082 let path = source.path().join("plugin");
2083 fs::write(&path, b"artifact").unwrap();
2084 let digest = format!("sha256:{}", hex::encode(Sha256::digest(b"artifact")));
2085
2086 let handle =
2087 ArtifactHandle::open_with_staging_root(&path, &digest, 8, staging.path()).unwrap();
2088
2089 assert_eq!(
2090 handle.path().parent().unwrap().parent().unwrap(),
2091 staging.path()
2092 );
2093 }
2094
2095 #[test]
2096 fn instance_resources_are_order_independent_and_immutable() {
2097 let left = InstanceResources::from_files([
2098 ("prompts/system.md".to_owned(), b"Build carefully.".to_vec()),
2099 ("rules.toml".to_owned(), b"turns = 4\n".to_vec()),
2100 ])
2101 .unwrap();
2102 let right = InstanceResources::from_files([
2103 ("rules.toml".to_owned(), b"turns = 4\n".to_vec()),
2104 ("prompts/system.md".to_owned(), b"Build carefully.".to_vec()),
2105 ])
2106 .unwrap();
2107
2108 assert_eq!(left.digest(), right.digest());
2109 assert_eq!(
2110 left.read_text("prompts/system.md").unwrap(),
2111 "Build carefully."
2112 );
2113 assert_eq!(left.file_count(), 2);
2114 assert_eq!(left.total_size(), 26);
2115 }
2116
2117 #[test]
2118 fn instance_resources_reject_escaping_and_duplicate_paths() {
2119 assert!(InstanceResources::from_files([("../secret".to_owned(), Vec::new())]).is_err());
2120 assert!(
2121 InstanceResources::from_files([
2122 ("rules.toml".to_owned(), Vec::new()),
2123 ("rules.toml".to_owned(), Vec::new()),
2124 ])
2125 .is_err()
2126 );
2127 }
2128}