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