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 let context = dependency.child_context(context)?;
1174 match binding
1175 .codec
1176 .open_host_stream(dependency, operation, request, context)
1177 .await?
1178 {
1179 Ok(stream) => {
1180 let stream_id = self.next_stream_id.get();
1181 let next =
1182 stream_id
1183 .checked_add(1)
1184 .ok_or(RuntimeFailure::ResourceExhausted {
1185 capability: JSON_HOST_IMPORTS_ABI_V2,
1186 operation: "stream-open".to_owned(),
1187 })?;
1188 self.next_stream_id.set(next);
1189 self.streams.borrow_mut().insert(stream_id, stream);
1190 Ok(Ok(stream_id))
1191 }
1192 Err(error) => Ok(Err(error)),
1193 }
1194 })
1195 }
1196
1197 pub fn send_stream(
1199 &self,
1200 stream_id: u64,
1201 message: Value,
1202 ) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>> {
1203 match self.stream(stream_id) {
1204 Ok(stream) => stream.send(message),
1205 Err(error) => Box::pin(futures::future::ready(Err(error))),
1206 }
1207 }
1208
1209 pub fn receive_stream(
1211 self: Rc<Self>,
1212 stream_id: u64,
1213 ) -> futures::future::LocalBoxFuture<'static, Result<JsonStreamItem, RuntimeFailure>> {
1214 Box::pin(async move {
1215 let stream = self.stream(stream_id)?;
1216 let item = stream.receive().await?;
1217 if matches!(item, JsonStreamItem::Terminal(_)) {
1218 self.streams.borrow_mut().remove(&stream_id);
1219 }
1220 Ok(item)
1221 })
1222 }
1223
1224 pub fn close_stream_send(
1226 &self,
1227 stream_id: u64,
1228 ) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>> {
1229 match self.stream(stream_id) {
1230 Ok(stream) => stream.close_send(),
1231 Err(error) => Box::pin(futures::future::ready(Err(error))),
1232 }
1233 }
1234
1235 pub fn cancel_stream(&self, stream_id: u64) -> Result<(), RuntimeFailure> {
1237 let stream = self
1238 .streams
1239 .borrow_mut()
1240 .remove(&stream_id)
1241 .ok_or_else(unknown_host_stream)?;
1242 stream.cancel();
1243 Ok(())
1244 }
1245
1246 pub fn deactivate(&self) {
1248 self.bindings.replace(None);
1249 for (_, stream) in std::mem::take(&mut *self.streams.borrow_mut()) {
1250 stream.cancel();
1251 }
1252 }
1253
1254 fn binding(&self, binding_id: u32) -> Result<JsonHostBinding, RuntimeFailure> {
1255 let bindings = self.bindings.borrow();
1256 let bindings = bindings.as_ref().ok_or(RuntimeFailure::AdmissionClosed)?;
1257 bindings
1258 .get(binding_id as usize)
1259 .cloned()
1260 .ok_or(RuntimeFailure::ProtocolViolation {
1261 capability: JSON_HOST_IMPORTS_ABI_V2,
1262 })
1263 }
1264
1265 fn stream(&self, stream_id: u64) -> Result<Rc<dyn JsonHostStreamSession>, RuntimeFailure> {
1266 self.streams
1267 .borrow()
1268 .get(&stream_id)
1269 .cloned()
1270 .ok_or_else(unknown_host_stream)
1271 }
1272}
1273
1274fn validate_host_binding(
1275 codec: &Rc<dyn JsonCapabilityCodec>,
1276 request: Option<&PluginDependencyHandle>,
1277 stream: Option<&PluginStreamDependencyHandle>,
1278 event: Option<&PluginEventDependencyHandle>,
1279) -> Result<(), RuntimeFailure> {
1280 for (capability, version) in request
1281 .map(|handle| (handle.capability_id(), handle.descriptor_version()))
1282 .into_iter()
1283 .chain(stream.map(|handle| (handle.capability_id(), handle.descriptor_version())))
1284 .chain(event.map(|handle| (handle.capability_id(), handle.descriptor_version())))
1285 {
1286 if capability != codec.capability_id() || version != codec.descriptor_version() {
1287 return Err(RuntimeFailure::ProtocolViolation {
1288 capability: codec.capability_id(),
1289 });
1290 }
1291 }
1292 Ok(())
1293}
1294
1295fn unknown_host_stream() -> RuntimeFailure {
1296 RuntimeFailure::ProtocolViolation {
1297 capability: JSON_HOST_IMPORTS_ABI_V2,
1298 }
1299}
1300
1301impl JsonStreamFrame {
1302 pub fn decode(
1304 encoded: &str,
1305 capability: &'static str,
1306 ) -> Result<JsonStreamItem, RuntimeFailure> {
1307 match serde_json::from_str(encoded)
1308 .map_err(|_| RuntimeFailure::ProtocolViolation { capability })?
1309 {
1310 Self::Message(value) => Ok(JsonStreamItem::Message(value)),
1311 Self::PeerHalfClosed => Ok(JsonStreamItem::PeerHalfClosed),
1312 Self::TerminalSuccess => Ok(JsonStreamItem::Terminal(Ok(()))),
1313 Self::TerminalError(value) => Ok(JsonStreamItem::Terminal(Err(value))),
1314 }
1315 }
1316}
1317
1318pub trait JsonStreamSessionTransport: std::fmt::Debug + 'static {
1320 fn send(
1321 self: Rc<Self>,
1322 message_json: String,
1323 ) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>>;
1324 fn receive(
1325 self: Rc<Self>,
1326 ) -> futures::future::LocalBoxFuture<'static, Result<JsonStreamItem, RuntimeFailure>>;
1327 fn close_send(
1328 self: Rc<Self>,
1329 ) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>>;
1330 fn cancel(&self);
1331}
1332
1333pub type JsonStreamOpenFuture = futures::future::LocalBoxFuture<
1335 'static,
1336 Result<Result<Rc<dyn JsonStreamSessionTransport>, Value>, RuntimeFailure>,
1337>;
1338
1339pub trait JsonStreamTransport: std::fmt::Debug + 'static {
1341 fn open(
1342 self: Rc<Self>,
1343 capability: String,
1344 operation: String,
1345 request_json: String,
1346 context: InvocationContext,
1347 ) -> JsonStreamOpenFuture;
1348}
1349
1350pub trait JsonEventTransport: std::fmt::Debug + 'static {
1352 fn publish(
1354 self: Rc<Self>,
1355 capability: String,
1356 operation: String,
1357 event_json: String,
1358 context: InvocationContext,
1359 ) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>>;
1360}
1361
1362pub fn json_request_endpoints<T: JsonRequestTransport>(
1364 transport: Rc<T>,
1365 codecs: Vec<Rc<dyn JsonCapabilityCodec>>,
1366) -> Vec<Rc<dyn NativeRequestEndpoint>> {
1367 let transport: Rc<dyn JsonRequestTransport> = transport;
1368 codecs
1369 .into_iter()
1370 .filter(|codec| !codec.request_operations().is_empty())
1371 .map(|codec| {
1372 Rc::new(JsonRequestEndpoint {
1373 transport: transport.clone(),
1374 codec,
1375 }) as Rc<dyn NativeRequestEndpoint>
1376 })
1377 .collect()
1378}
1379
1380pub fn json_stream_endpoints<T: JsonStreamTransport>(
1382 transport: Rc<T>,
1383 codecs: Vec<Rc<dyn JsonCapabilityCodec>>,
1384) -> Vec<Rc<dyn NativeStreamEndpoint>> {
1385 let transport: Rc<dyn JsonStreamTransport> = transport;
1386 codecs
1387 .into_iter()
1388 .filter(|codec| !codec.stream_operations().is_empty())
1389 .map(|codec| {
1390 Rc::new(JsonStreamEndpoint {
1391 transport: transport.clone(),
1392 codec,
1393 }) as Rc<dyn NativeStreamEndpoint>
1394 })
1395 .collect()
1396}
1397
1398pub fn json_event_endpoints<T: JsonEventTransport>(
1400 transport: Rc<T>,
1401 codecs: Vec<Rc<dyn JsonCapabilityCodec>>,
1402) -> Vec<Rc<dyn NativeEventEndpoint>> {
1403 let transport: Rc<dyn JsonEventTransport> = transport;
1404 codecs
1405 .into_iter()
1406 .filter(|codec| !codec.event_operations().is_empty())
1407 .map(|codec| {
1408 Rc::new(JsonEventEndpoint {
1409 transport: transport.clone(),
1410 codec,
1411 }) as Rc<dyn NativeEventEndpoint>
1412 })
1413 .collect()
1414}
1415
1416#[derive(Debug)]
1417struct JsonEventEndpoint {
1418 transport: Rc<dyn JsonEventTransport>,
1419 codec: Rc<dyn JsonCapabilityCodec>,
1420}
1421
1422impl NativeEventEndpoint for JsonEventEndpoint {
1423 fn capability_id(&self) -> &'static str {
1424 self.codec.capability_id()
1425 }
1426
1427 fn descriptor_version(&self) -> &'static str {
1428 self.codec.descriptor_version()
1429 }
1430
1431 fn operations(&self) -> &'static [&'static str] {
1432 self.codec.event_operations()
1433 }
1434
1435 fn owns_event_admission(&self) -> bool {
1436 true
1437 }
1438
1439 fn publish(
1440 &self,
1441 operation: &str,
1442 event: Box<dyn Any>,
1443 context: InvocationContext,
1444 ) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>> {
1445 if !self.codec.event_operations().contains(&operation) {
1446 return Box::pin(futures::future::ready(Err(unknown_operation(
1447 self.codec.capability_id(),
1448 operation,
1449 ))));
1450 }
1451 let event = self.codec.encode_event(operation, event.as_ref());
1452 let capability = self.codec.capability_id();
1453 let operation = operation.to_owned();
1454 let transport = self.transport.clone();
1455 Box::pin(async move {
1456 let event_json = serde_json::to_string(&event?)
1457 .map_err(|_| RuntimeFailure::ProtocolViolation { capability })?;
1458 transport
1459 .publish(capability.to_owned(), operation, event_json, context)
1460 .await
1461 })
1462 }
1463}
1464
1465#[derive(Debug)]
1466struct JsonStreamEndpoint {
1467 transport: Rc<dyn JsonStreamTransport>,
1468 codec: Rc<dyn JsonCapabilityCodec>,
1469}
1470
1471impl NativeStreamEndpoint for JsonStreamEndpoint {
1472 fn capability_id(&self) -> &'static str {
1473 self.codec.capability_id()
1474 }
1475 fn descriptor_version(&self) -> &'static str {
1476 self.codec.descriptor_version()
1477 }
1478 fn operations(&self) -> &'static [&'static str] {
1479 self.codec.stream_operations()
1480 }
1481
1482 fn open(
1483 &self,
1484 operation: &str,
1485 request: Box<dyn Any>,
1486 context: InvocationContext,
1487 ) -> futures::future::LocalBoxFuture<
1488 'static,
1489 Result<Result<Box<dyn NativeStreamSession>, Box<dyn Any>>, RuntimeFailure>,
1490 > {
1491 let transport = self.transport.clone();
1492 let codec = self.codec.clone();
1493 let operation = operation.to_owned();
1494 Box::pin(async move {
1495 if !codec.stream_operations().contains(&operation.as_str()) {
1496 return Err(unknown_operation(codec.capability_id(), &operation));
1497 }
1498 let request = codec.encode_stream_open(&operation, request.as_ref())?;
1499 let request_json =
1500 serde_json::to_string(&request).map_err(|_| RuntimeFailure::ProtocolViolation {
1501 capability: codec.capability_id(),
1502 })?;
1503 match transport
1504 .open(
1505 codec.capability_id().to_owned(),
1506 operation.clone(),
1507 request_json,
1508 context,
1509 )
1510 .await?
1511 {
1512 Ok(session) => Ok(Ok(Box::new(JsonStreamSession {
1513 session,
1514 codec,
1515 operation,
1516 }) as Box<dyn NativeStreamSession>)),
1517 Err(error) => codec.decode_stream_domain_error(&operation, error).map(Err),
1518 }
1519 })
1520 }
1521}
1522
1523#[derive(Debug)]
1524struct JsonStreamSession {
1525 session: Rc<dyn JsonStreamSessionTransport>,
1526 codec: Rc<dyn JsonCapabilityCodec>,
1527 operation: String,
1528}
1529
1530impl NativeStreamSession for JsonStreamSession {
1531 fn send(
1532 &self,
1533 message: Box<dyn Any>,
1534 ) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>> {
1535 let encoded = self
1536 .codec
1537 .encode_stream_message(&self.operation, message.as_ref())
1538 .and_then(|value| {
1539 serde_json::to_string(&value).map_err(|_| RuntimeFailure::ProtocolViolation {
1540 capability: self.codec.capability_id(),
1541 })
1542 });
1543 let session = self.session.clone();
1544 Box::pin(async move { session.send(encoded?).await })
1545 }
1546
1547 fn receive(
1548 &self,
1549 ) -> futures::future::LocalBoxFuture<'static, Result<NativeStreamItem, RuntimeFailure>> {
1550 let session = self.session.clone();
1551 let codec = self.codec.clone();
1552 let operation = self.operation.clone();
1553 Box::pin(async move {
1554 match session.receive().await? {
1555 JsonStreamItem::Message(value) => codec
1556 .decode_stream_message(&operation, value)
1557 .map(NativeStreamItem::Message),
1558 JsonStreamItem::PeerHalfClosed => Ok(NativeStreamItem::PeerHalfClosed),
1559 JsonStreamItem::Terminal(Ok(())) => Ok(NativeStreamItem::Terminal(Ok(()))),
1560 JsonStreamItem::Terminal(Err(value)) => codec
1561 .decode_stream_domain_error(&operation, value)
1562 .map(|error| NativeStreamItem::Terminal(Err(error))),
1563 }
1564 })
1565 }
1566
1567 fn close_send(&self) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>> {
1568 self.session.clone().close_send()
1569 }
1570
1571 fn cancel(&self) {
1572 self.session.cancel();
1573 }
1574}
1575
1576#[derive(Debug)]
1577struct JsonRequestEndpoint {
1578 transport: Rc<dyn JsonRequestTransport>,
1579 codec: Rc<dyn JsonCapabilityCodec>,
1580}
1581
1582impl NativeRequestEndpoint for JsonRequestEndpoint {
1583 fn capability_id(&self) -> &'static str {
1584 self.codec.capability_id()
1585 }
1586
1587 fn descriptor_version(&self) -> &'static str {
1588 self.codec.descriptor_version()
1589 }
1590
1591 fn operations(&self) -> &'static [&'static str] {
1592 self.codec.request_operations()
1593 }
1594
1595 fn invoke(
1596 &self,
1597 operation: &str,
1598 request: Box<dyn Any>,
1599 context: InvocationContext,
1600 ) -> futures::future::LocalBoxFuture<
1601 'static,
1602 Result<Result<Box<dyn Any>, Box<dyn Any>>, RuntimeFailure>,
1603 > {
1604 let transport = self.transport.clone();
1605 let codec = self.codec.clone();
1606 let operation = operation.to_owned();
1607 Box::pin(async move {
1608 if !codec.request_operations().contains(&operation.as_str()) {
1609 return Err(RuntimeFailure::UnknownOperation {
1610 capability: codec.capability_id(),
1611 operation,
1612 });
1613 }
1614 let request = codec.encode_request(&operation, request.as_ref())?;
1615 let request =
1616 serde_json::to_string(&request).map_err(|_| RuntimeFailure::ProtocolViolation {
1617 capability: codec.capability_id(),
1618 })?;
1619 match transport
1620 .invoke(
1621 codec.capability_id().to_owned(),
1622 operation.clone(),
1623 request,
1624 context,
1625 )
1626 .await?
1627 {
1628 JsonInvocationOutcome::Success(value) => {
1629 codec.decode_response(&operation, value).map(Ok)
1630 }
1631 JsonInvocationOutcome::DomainError(value) => {
1632 codec.decode_domain_error(&operation, value).map(Err)
1633 }
1634 }
1635 })
1636 }
1637}
1638
1639pub fn codecs_for_instance(
1641 instance: &PluginInstancePlan,
1642 codecs: &BTreeMap<String, Rc<dyn JsonCapabilityCodec>>,
1643) -> Result<Vec<Rc<dyn JsonCapabilityCodec>>, RuntimeFailure> {
1644 let mut selected = Vec::with_capacity(instance.provided_capabilities().len());
1645 for descriptor in instance.provided_capabilities() {
1646 let codec = codecs.get(descriptor.capability_id()).ok_or_else(|| {
1647 RuntimeFailure::InvalidResolvedPlan {
1648 detail: format!(
1649 "no generated codec for Capability `{}`",
1650 descriptor.capability_id()
1651 ),
1652 }
1653 })?;
1654 let request_operations: Vec<_> = codec
1655 .request_operations()
1656 .iter()
1657 .map(|operation| (*operation).to_owned())
1658 .collect();
1659 let stream_operations: Vec<_> = codec
1660 .stream_operations()
1661 .iter()
1662 .map(|operation| (*operation).to_owned())
1663 .collect();
1664 let event_operations: Vec<_> = codec
1665 .event_operations()
1666 .iter()
1667 .map(|operation| (*operation).to_owned())
1668 .collect();
1669 let expected_request: Vec<_> = descriptor
1670 .request_operations()
1671 .into_iter()
1672 .map(str::to_owned)
1673 .collect();
1674 let expected_stream: Vec<_> = descriptor
1675 .stream_operations()
1676 .into_iter()
1677 .map(str::to_owned)
1678 .collect();
1679 let expected_event: Vec<_> = descriptor
1680 .event_operations()
1681 .into_iter()
1682 .map(str::to_owned)
1683 .collect();
1684 if codec.descriptor_version() != descriptor.descriptor_version()
1685 || request_operations != expected_request
1686 || stream_operations != expected_stream
1687 || event_operations != expected_event
1688 {
1689 return Err(RuntimeFailure::ProtocolViolation {
1690 capability: codec.capability_id(),
1691 });
1692 }
1693 selected.push(codec.clone());
1694 }
1695 Ok(selected)
1696}
1697
1698pub fn codecs_for_requirements(
1700 instance: &PluginInstancePlan,
1701 codecs: &BTreeMap<String, Rc<dyn JsonCapabilityCodec>>,
1702) -> Result<Vec<Rc<dyn JsonCapabilityCodec>>, RuntimeFailure> {
1703 let mut selected = BTreeMap::new();
1704 for requirement in instance.required_capabilities() {
1705 let codec = codecs.get(requirement.capability_id()).ok_or_else(|| {
1706 RuntimeFailure::InvalidResolvedPlan {
1707 detail: format!(
1708 "no generated guest import codec for Capability `{}`",
1709 requirement.capability_id()
1710 ),
1711 }
1712 })?;
1713 if codec.descriptor_version() != requirement.descriptor_version() {
1714 return Err(RuntimeFailure::ProtocolViolation {
1715 capability: codec.capability_id(),
1716 });
1717 }
1718 selected
1719 .entry(requirement.capability_id().to_owned())
1720 .or_insert_with(|| codec.clone());
1721 }
1722 Ok(selected.into_values().collect())
1723}
1724
1725pub fn prepare_request_app(
1727 plan: &ResolvedAppPlan,
1728 execution_class: &ExecutionClassId,
1729 generations: BTreeMap<String, PreparedNativePlugin>,
1730) -> Result<PreparedNativeApp, RuntimeFailure> {
1731 let selected_instances = plan
1732 .plugin_instances()
1733 .iter()
1734 .filter(|instance| instance.execution_class() == execution_class)
1735 .map(|instance| instance.instance_key().to_owned())
1736 .collect::<std::collections::BTreeSet<_>>();
1737 let mut endpoints = BTreeMap::new();
1738 let mut stream_endpoints = BTreeMap::new();
1739 for (instance_key, generation) in &generations {
1740 for endpoint in generation.endpoints() {
1741 let identity = (instance_key.clone(), endpoint.capability_id().to_owned());
1742 if endpoints.insert(identity, endpoint.clone()).is_some() {
1743 return Err(RuntimeFailure::InvalidResolvedPlan {
1744 detail: format!("duplicate request endpoint on Instance `{instance_key}`"),
1745 });
1746 }
1747 }
1748 for endpoint in generation.stream_endpoints() {
1749 let identity = (instance_key.clone(), endpoint.capability_id().to_owned());
1750 if stream_endpoints
1751 .insert(identity, endpoint.clone())
1752 .is_some()
1753 {
1754 return Err(RuntimeFailure::InvalidResolvedPlan {
1755 detail: format!("duplicate stream endpoint on Instance `{instance_key}`"),
1756 });
1757 }
1758 }
1759 }
1760 for instance in plan
1761 .plugin_instances()
1762 .iter()
1763 .filter(|instance| selected_instances.contains(instance.instance_key()))
1764 {
1765 if !generations.contains_key(instance.instance_key()) {
1766 return Err(RuntimeFailure::InvalidResolvedPlan {
1767 detail: format!("Adapter omitted Instance `{}`", instance.instance_key()),
1768 });
1769 }
1770 }
1771 let mut bindings = Vec::new();
1772 let mut stream_bindings = Vec::new();
1773 for binding in plan.capability_bindings() {
1774 let key = (
1775 binding.provider_instance().to_owned(),
1776 binding.capability_id().to_owned(),
1777 );
1778 let request_endpoint = endpoints.get(&key);
1779 let stream_endpoint = stream_endpoints.get(&key);
1780 if let Some(endpoint) = request_endpoint {
1781 bindings.push(
1782 PreparedBinding::new(
1783 binding.consumer_instance(),
1784 binding.provider_instance(),
1785 endpoint.clone(),
1786 )
1787 .with_requirement_id(binding.requirement_id()),
1788 );
1789 }
1790 if let Some(endpoint) = stream_endpoint {
1791 stream_bindings.push(
1792 PreparedStreamBinding::new(
1793 binding.consumer_instance(),
1794 binding.provider_instance(),
1795 endpoint.clone(),
1796 )
1797 .with_requirement_id(binding.requirement_id()),
1798 );
1799 }
1800 if request_endpoint.is_none()
1801 && stream_endpoint.is_none()
1802 && selected_instances.contains(binding.provider_instance())
1803 {
1804 return Err(RuntimeFailure::InvalidResolvedPlan {
1805 detail: format!(
1806 "Adapter omitted Capability `{}` endpoint for Instance `{}`",
1807 binding.capability_id(),
1808 binding.provider_instance()
1809 ),
1810 });
1811 }
1812 }
1813 Ok(PreparedNativeApp::new(bindings, generations).with_stream_bindings(stream_bindings))
1814}
1815
1816pub fn require_operation(
1818 codecs: &BTreeMap<String, Rc<dyn JsonCapabilityCodec>>,
1819 capability_id: &str,
1820 operation: &str,
1821) -> Result<Rc<dyn JsonCapabilityCodec>, RuntimeFailure> {
1822 let codec =
1823 codecs
1824 .get(capability_id)
1825 .cloned()
1826 .ok_or_else(|| RuntimeFailure::InvalidResolvedPlan {
1827 detail: format!("no generated codec for Capability `{capability_id}`"),
1828 })?;
1829 if !codec.request_operations().contains(&operation) {
1830 return Err(RuntimeFailure::UnknownOperation {
1831 capability: codec.capability_id(),
1832 operation: operation.to_owned(),
1833 });
1834 }
1835 Ok(codec)
1836}
1837
1838fn unknown_operation(capability: &'static str, operation: &str) -> RuntimeFailure {
1839 RuntimeFailure::UnknownOperation {
1840 capability,
1841 operation: operation.to_owned(),
1842 }
1843}
1844
1845fn validate_digest(digest: &str) -> Result<(), RuntimeFailure> {
1846 let valid = digest.strip_prefix("sha256:").is_some_and(|hex| {
1847 hex.len() == 64
1848 && hex
1849 .bytes()
1850 .all(|byte| byte.is_ascii_hexdigit() && !byte.is_ascii_uppercase())
1851 });
1852 if valid {
1853 Ok(())
1854 } else {
1855 Err(RuntimeFailure::InvalidResolvedPlan {
1856 detail: format!("invalid canonical SHA-256 digest `{digest}`"),
1857 })
1858 }
1859}
1860
1861fn invalid_artifact(path: &Path, error: impl std::fmt::Display) -> RuntimeFailure {
1862 RuntimeFailure::InvalidResolvedPlan {
1863 detail: format!("cannot read Artifact `{}`: {error}", path.display()),
1864 }
1865}
1866
1867fn validate_resource_path(path: &str) -> Result<(), RuntimeFailure> {
1868 if path.is_empty()
1869 || path.starts_with('/')
1870 || path.contains(['\\', '\0'])
1871 || path
1872 .split('/')
1873 .any(|segment| segment.is_empty() || matches!(segment, "." | ".."))
1874 {
1875 return Err(invalid_resources(format!(
1876 "invalid Plugin resource path `{path}`"
1877 )));
1878 }
1879 Ok(())
1880}
1881
1882fn invalid_resources(detail: impl Into<String>) -> RuntimeFailure {
1883 RuntimeFailure::InvalidResolvedPlan {
1884 detail: detail.into(),
1885 }
1886}
1887
1888#[cfg(test)]
1889mod tests {
1890 use std::{cell::RefCell, io::Write};
1891
1892 use super::*;
1893
1894 #[derive(Debug)]
1895 struct EventCodec;
1896
1897 impl JsonCapabilityCodec for EventCodec {
1898 fn capability_id(&self) -> &'static str {
1899 "example.events@1"
1900 }
1901
1902 fn descriptor_version(&self) -> &'static str {
1903 "1.0.0"
1904 }
1905
1906 fn request_operations(&self) -> &'static [&'static str] {
1907 &[]
1908 }
1909
1910 fn event_operations(&self) -> &'static [&'static str] {
1911 &["changed"]
1912 }
1913
1914 fn encode_request(
1915 &self,
1916 operation: &str,
1917 _request: &dyn Any,
1918 ) -> Result<Value, RuntimeFailure> {
1919 Err(unknown_operation(self.capability_id(), operation))
1920 }
1921
1922 fn decode_response(
1923 &self,
1924 operation: &str,
1925 _value: Value,
1926 ) -> Result<Box<dyn Any>, RuntimeFailure> {
1927 Err(unknown_operation(self.capability_id(), operation))
1928 }
1929
1930 fn decode_domain_error(
1931 &self,
1932 operation: &str,
1933 _value: Value,
1934 ) -> Result<Box<dyn Any>, RuntimeFailure> {
1935 Err(unknown_operation(self.capability_id(), operation))
1936 }
1937
1938 fn encode_event(&self, operation: &str, event: &dyn Any) -> Result<Value, RuntimeFailure> {
1939 if operation != "changed" {
1940 return Err(unknown_operation(self.capability_id(), operation));
1941 }
1942 event
1943 .downcast_ref::<String>()
1944 .map(|value| serde_json::json!({"value": value}))
1945 .ok_or(RuntimeFailure::ProtocolViolation {
1946 capability: self.capability_id(),
1947 })
1948 }
1949 }
1950
1951 #[derive(Debug, Default)]
1952 struct EventTransport {
1953 published: RefCell<Option<(String, String, String)>>,
1954 }
1955
1956 impl JsonEventTransport for EventTransport {
1957 fn publish(
1958 self: Rc<Self>,
1959 capability: String,
1960 operation: String,
1961 event_json: String,
1962 _context: InvocationContext,
1963 ) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>> {
1964 Box::pin(async move {
1965 self.published
1966 .replace(Some((capability, operation, event_json)));
1967 Ok(())
1968 })
1969 }
1970 }
1971
1972 #[test]
1973 fn event_endpoints_encode_typed_values_and_commit_transport_admission() {
1974 let transport = Rc::new(EventTransport::default());
1975 let endpoints = json_event_endpoints(
1976 transport.clone(),
1977 vec![Rc::new(EventCodec) as Rc<dyn JsonCapabilityCodec>],
1978 );
1979 assert_eq!(endpoints.len(), 1);
1980 assert!(endpoints[0].owns_event_admission());
1981 assert_eq!(endpoints[0].operations(), &["changed"]);
1982
1983 futures::executor::block_on(endpoints[0].publish(
1984 "changed",
1985 Box::new("ready".to_owned()),
1986 InvocationContext::new(1, None, lenso_kernel::CancellationToken::new()),
1987 ))
1988 .unwrap();
1989
1990 assert_eq!(
1991 transport.published.borrow().as_ref().unwrap(),
1992 &(
1993 "example.events@1".to_owned(),
1994 "changed".to_owned(),
1995 r#"{"value":"ready"}"#.to_owned(),
1996 )
1997 );
1998 }
1999
2000 #[test]
2001 fn artifact_handle_keeps_the_admitted_bytes_after_source_drift() {
2002 let mut file = tempfile::NamedTempFile::new().unwrap();
2003 file.write_all(b"first").unwrap();
2004 let digest = format!("sha256:{}", hex::encode(Sha256::digest(b"first")));
2005 let handle = ArtifactHandle::open(file.path(), &digest, 5).unwrap();
2006 file.as_file_mut().set_len(0).unwrap();
2007 file.write_all(b"other").unwrap();
2008 assert_eq!(handle.read_verified().unwrap(), b"first");
2009 assert_ne!(handle.path(), file.path());
2010 assert_eq!(fs::read(handle.path()).unwrap(), b"first");
2011 }
2012
2013 #[test]
2014 fn artifact_snapshot_survives_source_parent_rename_and_replacement() {
2015 let workspace = tempfile::tempdir().unwrap();
2016 let selected = workspace.path().join("selected");
2017 fs::create_dir(&selected).unwrap();
2018 let source = selected.join("plugin");
2019 fs::write(&source, b"admitted").unwrap();
2020 let digest = format!("sha256:{}", hex::encode(Sha256::digest(b"admitted")));
2021
2022 let handle = ArtifactHandle::open(&source, &digest, 8).unwrap();
2023 assert!(!handle.path().starts_with(&selected));
2024 fs::rename(&selected, workspace.path().join("replaced")).unwrap();
2025 fs::create_dir(&selected).unwrap();
2026 fs::write(selected.join("plugin"), b"attacker").unwrap();
2027
2028 assert_eq!(fs::read(handle.path()).unwrap(), b"admitted");
2029 assert_eq!(handle.read_verified().unwrap(), b"admitted");
2030 }
2031
2032 #[test]
2033 fn artifact_admission_streams_large_content_into_one_stable_snapshot() {
2034 let bytes = vec![0x5a; 4 * 1024 * 1024 + 17];
2035 let mut file = tempfile::NamedTempFile::new().unwrap();
2036 file.write_all(&bytes).unwrap();
2037 let digest = format!("sha256:{}", hex::encode(Sha256::digest(&bytes)));
2038
2039 let handle = ArtifactHandle::open(file.path(), &digest, bytes.len() as u64).unwrap();
2040
2041 assert_eq!(handle.read_verified().unwrap(), bytes);
2042 }
2043
2044 #[test]
2047 #[ignore = "large Artifact admission benchmark; run explicitly"]
2048 fn artifact_admission_streaming_benchmark() {
2049 const BLOCK_BYTES: usize = 64 * 1024;
2050 let block = vec![0x5a; BLOCK_BYTES];
2051 for mebibytes in [4_usize, 64, 256] {
2052 let directory = tempfile::tempdir().unwrap();
2053 let path = directory.path().join("artifact");
2054 let mut source = fs::File::create(&path).unwrap();
2055 let mut hasher = Sha256::new();
2056 let blocks = mebibytes * 1024 * 1024 / BLOCK_BYTES;
2057 for _ in 0..blocks {
2058 source.write_all(&block).unwrap();
2059 hasher.update(&block);
2060 }
2061 drop(source);
2062 let size = u64::try_from(mebibytes * 1024 * 1024).unwrap();
2063 let digest = format!("sha256:{}", hex::encode(hasher.finalize()));
2064
2065 let started = std::time::Instant::now();
2066 let handle = ArtifactHandle::open(&path, &digest, size).unwrap();
2067 let elapsed = started.elapsed();
2068
2069 assert_eq!(handle.size(), size);
2070 println!(
2071 "{{\"mebibytes\":{mebibytes},\"elapsed_ms\":{:.3},\"mib_per_second\":{:.3}}}",
2072 elapsed.as_secs_f64() * 1_000.0,
2073 f64::from(u32::try_from(mebibytes).unwrap()) / elapsed.as_secs_f64()
2074 );
2075 drop(handle);
2076 }
2077 }
2078
2079 #[test]
2080 fn artifact_admission_honors_an_explicit_host_staging_root() {
2081 let source = tempfile::tempdir().unwrap();
2082 let staging = tempfile::tempdir().unwrap();
2083 let path = source.path().join("plugin");
2084 fs::write(&path, b"artifact").unwrap();
2085 let digest = format!("sha256:{}", hex::encode(Sha256::digest(b"artifact")));
2086
2087 let handle =
2088 ArtifactHandle::open_with_staging_root(&path, &digest, 8, staging.path()).unwrap();
2089
2090 assert_eq!(
2091 handle.path().parent().unwrap().parent().unwrap(),
2092 staging.path()
2093 );
2094 }
2095
2096 #[test]
2097 fn instance_resources_are_order_independent_and_immutable() {
2098 let left = InstanceResources::from_files([
2099 ("prompts/system.md".to_owned(), b"Build carefully.".to_vec()),
2100 ("rules.toml".to_owned(), b"turns = 4\n".to_vec()),
2101 ])
2102 .unwrap();
2103 let right = InstanceResources::from_files([
2104 ("rules.toml".to_owned(), b"turns = 4\n".to_vec()),
2105 ("prompts/system.md".to_owned(), b"Build carefully.".to_vec()),
2106 ])
2107 .unwrap();
2108
2109 assert_eq!(left.digest(), right.digest());
2110 assert_eq!(
2111 left.read_text("prompts/system.md").unwrap(),
2112 "Build carefully."
2113 );
2114 assert_eq!(left.file_count(), 2);
2115 assert_eq!(left.total_size(), 26);
2116 }
2117
2118 #[test]
2119 fn instance_resources_reject_escaping_and_duplicate_paths() {
2120 assert!(InstanceResources::from_files([("../secret".to_owned(), Vec::new())]).is_err());
2121 assert!(
2122 InstanceResources::from_files([
2123 ("rules.toml".to_owned(), Vec::new()),
2124 ("rules.toml".to_owned(), Vec::new()),
2125 ])
2126 .is_err()
2127 );
2128 }
2129}