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