use std::{
any::Any,
collections::BTreeMap,
fs,
path::{Path, PathBuf},
rc::Rc,
};
use lenso_app_plan::{
CapabilityCardinality, ExecutionClassId, PluginInstancePlan, ResolvedAppPlan,
};
use lenso_kernel::{
InvocationContext, NativeRequestEndpoint, NativeStream, NativeStreamEndpoint, NativeStreamItem,
NativeStreamSession, PluginDependencies, PluginDependencyHandle, PluginStreamDependencyHandle,
PreparedBinding, PreparedNativeApp, PreparedNativePlugin, PreparedStreamBinding,
RuntimeFailure, StreamCapability, StreamEvent,
};
use serde::{Deserialize, Serialize};
use serde_json::Value;
use sha2::{Digest, Sha256};
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct ArtifactHandle {
path: PathBuf,
digest: String,
size: u64,
}
impl ArtifactHandle {
pub fn open(
path: impl Into<PathBuf>,
expected_digest: &str,
expected_size: u64,
) -> Result<Self, RuntimeFailure> {
validate_digest(expected_digest)?;
let path = path.into();
let metadata =
fs::symlink_metadata(&path).map_err(|error| invalid_artifact(&path, error))?;
if !metadata.file_type().is_file() || metadata.file_type().is_symlink() {
return Err(RuntimeFailure::InvalidResolvedPlan {
detail: format!("Artifact `{}` is not a regular file", path.display()),
});
}
if metadata.len() != expected_size {
return Err(RuntimeFailure::InvalidResolvedPlan {
detail: format!(
"Artifact `{}` size mismatch: expected {expected_size}, got {}",
path.display(),
metadata.len()
),
});
}
let bytes = fs::read(&path).map_err(|error| invalid_artifact(&path, error))?;
let actual_digest = format!("sha256:{}", hex::encode(Sha256::digest(&bytes)));
if actual_digest != expected_digest {
return Err(RuntimeFailure::InvalidResolvedPlan {
detail: format!("Artifact `{}` digest mismatch", path.display()),
});
}
Ok(Self {
path,
digest: actual_digest,
size: metadata.len(),
})
}
pub fn path(&self) -> &Path {
&self.path
}
pub fn digest(&self) -> &str {
&self.digest
}
pub const fn size(&self) -> u64 {
self.size
}
pub fn read_verified(&self) -> Result<Vec<u8>, RuntimeFailure> {
let verified = Self::open(&self.path, &self.digest, self.size)?;
fs::read(verified.path).map_err(|error| invalid_artifact(&self.path, error))
}
}
#[derive(Clone, Debug, Default)]
pub struct ArtifactCatalog(BTreeMap<String, ArtifactHandle>);
impl ArtifactCatalog {
pub fn new() -> Self {
Self::default()
}
pub fn with_artifact(
mut self,
instance_key: impl Into<String>,
artifact: ArtifactHandle,
) -> Result<Self, RuntimeFailure> {
let instance_key = instance_key.into();
if self.0.insert(instance_key.clone(), artifact).is_some() {
return Err(RuntimeFailure::InvalidResolvedPlan {
detail: format!("duplicate Artifact authority for Instance `{instance_key}`"),
});
}
Ok(self)
}
pub fn require(&self, instance_key: &str) -> Result<&ArtifactHandle, RuntimeFailure> {
self.0
.get(instance_key)
.ok_or_else(|| RuntimeFailure::InvalidResolvedPlan {
detail: format!("no admitted Artifact for Instance `{instance_key}`"),
})
}
}
pub trait JsonCapabilityCodec: std::fmt::Debug + 'static {
fn capability_id(&self) -> &'static str;
fn descriptor_version(&self) -> &'static str;
fn request_operations(&self) -> &'static [&'static str];
fn stream_operations(&self) -> &'static [&'static str] {
&[]
}
fn encode_request(&self, operation: &str, request: &dyn Any) -> Result<Value, RuntimeFailure>;
fn decode_response(
&self,
operation: &str,
value: Value,
) -> Result<Box<dyn Any>, RuntimeFailure>;
fn decode_domain_error(
&self,
operation: &str,
value: Value,
) -> Result<Box<dyn Any>, RuntimeFailure>;
fn encode_stream_open(
&self,
operation: &str,
request: &dyn Any,
) -> Result<Value, RuntimeFailure> {
let _ = request;
Err(unknown_operation(self.capability_id(), operation))
}
fn encode_stream_message(
&self,
operation: &str,
message: &dyn Any,
) -> Result<Value, RuntimeFailure> {
let _ = message;
Err(unknown_operation(self.capability_id(), operation))
}
fn decode_stream_message(
&self,
operation: &str,
value: Value,
) -> Result<Box<dyn Any>, RuntimeFailure> {
let _ = value;
Err(unknown_operation(self.capability_id(), operation))
}
fn decode_stream_domain_error(
&self,
operation: &str,
value: Value,
) -> Result<Box<dyn Any>, RuntimeFailure> {
let _ = value;
Err(unknown_operation(self.capability_id(), operation))
}
fn invoke_host_request(
&self,
dependency: PluginDependencyHandle,
operation: String,
request: Value,
context: InvocationContext,
) -> JsonHostRequestFuture {
let _ = (dependency, request, context);
Box::pin(futures::future::ready(Err(unknown_operation(
self.capability_id(),
&operation,
))))
}
fn open_host_stream(
&self,
dependency: PluginStreamDependencyHandle,
operation: String,
request: Value,
context: InvocationContext,
) -> JsonHostStreamOpenFuture {
let _ = (dependency, request, context);
Box::pin(futures::future::ready(Err(unknown_operation(
self.capability_id(),
&operation,
))))
}
}
#[derive(Debug)]
pub enum JsonInvocationOutcome {
Success(Value),
DomainError(Value),
}
pub fn json_runtime_failure(error: &RuntimeFailure) -> Value {
match error {
RuntimeFailure::Unavailable { capability } => serde_json::json!({
"kind": "unavailable",
"capability": capability,
}),
RuntimeFailure::UnknownOperation {
capability,
operation,
} => serde_json::json!({
"kind": "unknown_operation",
"capability": capability,
"operation": operation,
}),
RuntimeFailure::AmbiguousBinding {
capability,
providers,
} => serde_json::json!({
"kind": "ambiguous_binding",
"capability": capability,
"providers": providers,
}),
RuntimeFailure::ProtocolViolation { capability } => serde_json::json!({
"kind": "protocol_violation",
"capability": capability,
}),
RuntimeFailure::AdmissionClosed => serde_json::json!({ "kind": "admission_closed" }),
RuntimeFailure::ResourceExhausted {
capability,
operation,
} => serde_json::json!({
"kind": "resource_exhausted",
"capability": capability,
"operation": operation,
}),
RuntimeFailure::DeadlineExceeded { request_id } => serde_json::json!({
"kind": "deadline_exceeded",
"request_id": request_id.to_string(),
}),
RuntimeFailure::Cancelled { request_id } => serde_json::json!({
"kind": "cancelled",
"request_id": request_id.to_string(),
}),
RuntimeFailure::MissingPluginFactory { .. }
| RuntimeFailure::UnavailableExecutionClass { .. }
| RuntimeFailure::InvalidResolvedPlan { .. }
| RuntimeFailure::Internal { .. }
| RuntimeFailure::PluginFailure { .. }
| RuntimeFailure::PluginRestartExhausted { .. } => {
serde_json::json!({ "kind": "internal" })
}
}
}
pub fn json_host_invocation_envelope(
outcome: Result<JsonInvocationOutcome, RuntimeFailure>,
) -> Value {
match outcome {
Ok(JsonInvocationOutcome::Success(value)) => serde_json::json!({ "ok": value }),
Ok(JsonInvocationOutcome::DomainError(value)) => serde_json::json!({ "error": value }),
Err(error) => serde_json::json!({ "runtime": json_runtime_failure(&error) }),
}
}
pub type JsonHostRequestFuture =
futures::future::LocalBoxFuture<'static, Result<JsonInvocationOutcome, RuntimeFailure>>;
pub trait JsonHostStreamSession: std::fmt::Debug + 'static {
fn send(
self: Rc<Self>,
message: Value,
) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>>;
fn receive(
self: Rc<Self>,
) -> futures::future::LocalBoxFuture<'static, Result<JsonStreamItem, RuntimeFailure>>;
fn close_send(
self: Rc<Self>,
) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>>;
fn cancel(&self);
}
pub type JsonHostStreamOpenFuture = futures::future::LocalBoxFuture<
'static,
Result<Result<Rc<dyn JsonHostStreamSession>, Value>, RuntimeFailure>,
>;
type DecodeStreamMessage<C> =
Rc<dyn Fn(Value) -> Result<<C as StreamCapability>::Message, RuntimeFailure>>;
type EncodeStreamMessage<C> =
Rc<dyn Fn(<C as StreamCapability>::Message) -> Result<Value, RuntimeFailure>>;
type EncodeStreamError<C> =
Rc<dyn Fn(<C as StreamCapability>::DomainError) -> Result<Value, RuntimeFailure>>;
pub fn json_host_stream<C: StreamCapability>(
stream: NativeStream<C>,
decode_message: impl Fn(Value) -> Result<C::Message, RuntimeFailure> + 'static,
encode_message: impl Fn(C::Message) -> Result<Value, RuntimeFailure> + 'static,
encode_error: impl Fn(C::DomainError) -> Result<Value, RuntimeFailure> + 'static,
) -> Rc<dyn JsonHostStreamSession> {
Rc::new(TypedJsonHostStream {
stream: Rc::new(stream),
decode_message: Rc::new(decode_message),
encode_message: Rc::new(encode_message),
encode_error: Rc::new(encode_error),
})
}
struct TypedJsonHostStream<C: StreamCapability> {
stream: Rc<NativeStream<C>>,
decode_message: DecodeStreamMessage<C>,
encode_message: EncodeStreamMessage<C>,
encode_error: EncodeStreamError<C>,
}
impl<C: StreamCapability> std::fmt::Debug for TypedJsonHostStream<C> {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter
.debug_struct("TypedJsonHostStream")
.field("capability", &C::ID)
.finish_non_exhaustive()
}
}
impl<C: StreamCapability> JsonHostStreamSession for TypedJsonHostStream<C> {
fn send(
self: Rc<Self>,
message: Value,
) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>> {
Box::pin(async move {
let message = (self.decode_message)(message)?;
self.stream.send(message).await
})
}
fn receive(
self: Rc<Self>,
) -> futures::future::LocalBoxFuture<'static, Result<JsonStreamItem, RuntimeFailure>> {
Box::pin(async move {
match self.stream.receive().await? {
StreamEvent::Message(message) => {
(self.encode_message)(message).map(JsonStreamItem::Message)
}
StreamEvent::PeerHalfClosed => Ok(JsonStreamItem::PeerHalfClosed),
StreamEvent::Terminal(Ok(())) => Ok(JsonStreamItem::Terminal(Ok(()))),
StreamEvent::Terminal(Err(error)) => {
(self.encode_error)(error).map(|error| JsonStreamItem::Terminal(Err(error)))
}
}
})
}
fn close_send(
self: Rc<Self>,
) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>> {
Box::pin(async move { self.stream.close_send().await })
}
fn cancel(&self) {
self.stream.cancel();
}
}
pub const JSON_REQUEST_ABI_V1: &str = "lenso.json-request@1";
pub const JSON_INTERACTIONS_ABI_V1: &str = "lenso.json-interactions@1";
pub const JSON_HOST_IMPORTS_ABI_V1: &str = "lenso.json-host-imports@1";
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(deny_unknown_fields)]
pub struct JsonPluginDescriptor {
pub abi: String,
pub capabilities: Vec<JsonCapabilityDescriptor>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub required_capabilities: Vec<JsonRequiredCapabilityDescriptor>,
}
#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
#[serde(deny_unknown_fields)]
pub struct JsonCapabilityDescriptor {
pub capability_id: String,
pub descriptor_version: String,
pub request_operations: Vec<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub stream_operations: Vec<String>,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(deny_unknown_fields)]
pub struct JsonRequiredCapabilityDescriptor {
pub capability_id: String,
pub descriptor_version: String,
pub cardinality: CapabilityCardinality,
}
pub fn expected_json_plugin_descriptor(
instance: &PluginInstancePlan,
) -> Result<JsonPluginDescriptor, RuntimeFailure> {
let mut capabilities = Vec::with_capacity(instance.provided_capabilities().len());
for descriptor in instance.provided_capabilities() {
if !descriptor.event_operations().is_empty() {
return Err(RuntimeFailure::InvalidResolvedPlan {
detail: format!(
"Execution class `{}` does not support Event endpoints",
instance.execution_class()
),
});
}
capabilities.push(JsonCapabilityDescriptor {
capability_id: descriptor.capability_id().to_owned(),
descriptor_version: descriptor.descriptor_version().to_owned(),
request_operations: descriptor
.request_operations()
.into_iter()
.map(str::to_owned)
.collect(),
stream_operations: descriptor
.stream_operations()
.into_iter()
.map(str::to_owned)
.collect(),
});
}
capabilities.sort();
if capabilities
.windows(2)
.any(|pair| pair[0].capability_id == pair[1].capability_id)
{
return Err(RuntimeFailure::InvalidResolvedPlan {
detail: format!(
"Instance `{}` declares a duplicate Capability",
instance.instance_key()
),
});
}
let mut required_capabilities = instance
.required_capabilities()
.iter()
.map(|requirement| JsonRequiredCapabilityDescriptor {
capability_id: requirement.capability_id().to_owned(),
descriptor_version: requirement.descriptor_version().to_owned(),
cardinality: requirement.cardinality(),
})
.collect::<Vec<_>>();
sort_required_capabilities(&mut required_capabilities);
Ok(JsonPluginDescriptor {
abi: if !required_capabilities.is_empty() {
JSON_HOST_IMPORTS_ABI_V1
} else if capabilities
.iter()
.any(|capability| !capability.stream_operations.is_empty())
{
JSON_INTERACTIONS_ABI_V1
} else {
JSON_REQUEST_ABI_V1
}
.to_owned(),
capabilities,
required_capabilities,
})
}
pub fn validate_json_plugin_descriptor(
instance: &PluginInstancePlan,
encoded: &str,
) -> Result<(), RuntimeFailure> {
let mut actual = serde_json::from_str::<JsonPluginDescriptor>(encoded).map_err(|_| {
RuntimeFailure::ProtocolViolation {
capability: "lenso.json-request@1",
}
})?;
actual.capabilities.sort();
sort_required_capabilities(&mut actual.required_capabilities);
let expected = expected_json_plugin_descriptor(instance)?;
if actual != expected {
return Err(RuntimeFailure::InvalidResolvedPlan {
detail: format!(
"guest descriptor does not match resolved Instance `{}`",
instance.instance_key()
),
});
}
Ok(())
}
fn sort_required_capabilities(requirements: &mut [JsonRequiredCapabilityDescriptor]) {
requirements.sort_by(|left, right| {
(
&left.capability_id,
&left.descriptor_version,
cardinality_order(left.cardinality),
)
.cmp(&(
&right.capability_id,
&right.descriptor_version,
cardinality_order(right.cardinality),
))
});
}
const fn cardinality_order(cardinality: CapabilityCardinality) -> u8 {
match cardinality {
CapabilityCardinality::One => 0,
CapabilityCardinality::Optional => 1,
CapabilityCardinality::Many => 2,
}
}
pub trait JsonRequestTransport: std::fmt::Debug + 'static {
fn invoke(
self: Rc<Self>,
capability: String,
operation: String,
request_json: String,
context: InvocationContext,
) -> futures::future::LocalBoxFuture<'static, Result<JsonInvocationOutcome, RuntimeFailure>>;
}
#[derive(Debug)]
pub enum JsonStreamItem {
Message(Value),
PeerHalfClosed,
Terminal(Result<(), Value>),
}
#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
#[serde(
tag = "kind",
content = "value",
rename_all = "kebab-case",
deny_unknown_fields
)]
pub enum JsonStreamFrame {
Message(Value),
PeerHalfClosed,
TerminalSuccess,
TerminalError(Value),
}
#[derive(Clone, Debug, Eq, PartialEq, Serialize)]
pub struct JsonHostBindingDescriptor {
pub binding_id: u32,
pub provider_instance: String,
pub capability_id: String,
pub descriptor_version: String,
pub request_operations: Vec<String>,
pub stream_operations: Vec<String>,
}
#[derive(Clone)]
struct JsonHostBinding {
descriptor: JsonHostBindingDescriptor,
codec: Rc<dyn JsonCapabilityCodec>,
request: Option<PluginDependencyHandle>,
stream: Option<PluginStreamDependencyHandle>,
}
impl std::fmt::Debug for JsonHostBinding {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter
.debug_struct("JsonHostBinding")
.field("descriptor", &self.descriptor)
.finish_non_exhaustive()
}
}
#[derive(Debug)]
pub struct JsonHostImports {
codecs: BTreeMap<String, Rc<dyn JsonCapabilityCodec>>,
bindings: std::cell::RefCell<Option<Vec<JsonHostBinding>>>,
streams: std::cell::RefCell<BTreeMap<u64, Rc<dyn JsonHostStreamSession>>>,
next_stream_id: std::cell::Cell<u64>,
max_streams: usize,
}
impl JsonHostImports {
pub fn new(
codecs: Vec<Rc<dyn JsonCapabilityCodec>>,
max_streams: usize,
) -> Result<Self, RuntimeFailure> {
let mut by_capability = BTreeMap::new();
for codec in codecs {
let capability = codec.capability_id().to_owned();
if by_capability.insert(capability.clone(), codec).is_some() {
return Err(RuntimeFailure::InvalidResolvedPlan {
detail: format!("duplicate guest import codec for Capability `{capability}`"),
});
}
}
Ok(Self {
codecs: by_capability,
bindings: std::cell::RefCell::new(None),
streams: std::cell::RefCell::new(BTreeMap::new()),
next_stream_id: std::cell::Cell::new(1),
max_streams,
})
}
pub fn activate(&self, dependencies: &PluginDependencies) -> Result<(), RuntimeFailure> {
if self.bindings.borrow().is_some() {
return Err(RuntimeFailure::Internal {
detail: "guest Capability imports were activated twice".to_owned(),
});
}
let mut bindings = Vec::with_capacity(dependencies.len());
for (index, dependency) in dependencies.bindings().iter().enumerate() {
let codec = self
.codecs
.get(dependency.capability_id())
.cloned()
.ok_or_else(|| RuntimeFailure::InvalidResolvedPlan {
detail: format!(
"no generated guest import codec for Capability `{}`",
dependency.capability_id()
),
})?;
let request = dependency.handle();
let stream = dependency.stream_handle();
validate_host_binding(&codec, request.as_ref(), stream.as_ref())?;
let binding_id =
u32::try_from(index).map_err(|_| RuntimeFailure::InvalidResolvedPlan {
detail: "guest import binding table exceeds u32 identity space".to_owned(),
})?;
bindings.push(JsonHostBinding {
descriptor: JsonHostBindingDescriptor {
binding_id,
provider_instance: dependency.provider_instance().to_owned(),
capability_id: dependency.capability_id().to_owned(),
descriptor_version: codec.descriptor_version().to_owned(),
request_operations: request.as_ref().map_or_else(Vec::new, |handle| {
handle
.operations()
.iter()
.map(|item| (*item).to_owned())
.collect()
}),
stream_operations: stream.as_ref().map_or_else(Vec::new, |handle| {
handle
.operations()
.iter()
.map(|item| (*item).to_owned())
.collect()
}),
},
codec,
request,
stream,
});
}
self.bindings.replace(Some(bindings));
Ok(())
}
pub fn descriptors(&self) -> Result<Vec<JsonHostBindingDescriptor>, RuntimeFailure> {
self.bindings
.borrow()
.as_ref()
.map(|bindings| {
bindings
.iter()
.map(|binding| binding.descriptor.clone())
.collect()
})
.ok_or(RuntimeFailure::AdmissionClosed)
}
pub fn invoke(
&self,
binding_id: u32,
operation: String,
request: Value,
context: InvocationContext,
) -> JsonHostRequestFuture {
let binding = match self.binding(binding_id) {
Ok(binding) => binding,
Err(error) => return Box::pin(futures::future::ready(Err(error))),
};
let Some(dependency) = binding.request else {
return Box::pin(futures::future::ready(Err(
RuntimeFailure::UnknownOperation {
capability: binding.codec.capability_id(),
operation,
},
)));
};
binding
.codec
.invoke_host_request(dependency, operation, request, context)
}
pub fn open_stream(
self: Rc<Self>,
binding_id: u32,
operation: String,
request: Value,
context: InvocationContext,
) -> futures::future::LocalBoxFuture<'static, Result<Result<u64, Value>, RuntimeFailure>> {
Box::pin(async move {
if self.streams.borrow().len() >= self.max_streams {
return Err(RuntimeFailure::ResourceExhausted {
capability: "lenso.json-host-imports@1",
operation: "stream-open".to_owned(),
});
}
let binding = self.binding(binding_id)?;
let dependency = binding
.stream
.ok_or_else(|| RuntimeFailure::UnknownOperation {
capability: binding.codec.capability_id(),
operation: operation.clone(),
})?;
match binding
.codec
.open_host_stream(dependency, operation, request, context)
.await?
{
Ok(stream) => {
let stream_id = self.next_stream_id.get();
let next =
stream_id
.checked_add(1)
.ok_or(RuntimeFailure::ResourceExhausted {
capability: "lenso.json-host-imports@1",
operation: "stream-open".to_owned(),
})?;
self.next_stream_id.set(next);
self.streams.borrow_mut().insert(stream_id, stream);
Ok(Ok(stream_id))
}
Err(error) => Ok(Err(error)),
}
})
}
pub fn send_stream(
&self,
stream_id: u64,
message: Value,
) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>> {
match self.stream(stream_id) {
Ok(stream) => stream.send(message),
Err(error) => Box::pin(futures::future::ready(Err(error))),
}
}
pub fn receive_stream(
self: Rc<Self>,
stream_id: u64,
) -> futures::future::LocalBoxFuture<'static, Result<JsonStreamItem, RuntimeFailure>> {
Box::pin(async move {
let stream = self.stream(stream_id)?;
let item = stream.receive().await?;
if matches!(item, JsonStreamItem::Terminal(_)) {
self.streams.borrow_mut().remove(&stream_id);
}
Ok(item)
})
}
pub fn close_stream_send(
&self,
stream_id: u64,
) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>> {
match self.stream(stream_id) {
Ok(stream) => stream.close_send(),
Err(error) => Box::pin(futures::future::ready(Err(error))),
}
}
pub fn cancel_stream(&self, stream_id: u64) -> Result<(), RuntimeFailure> {
let stream = self
.streams
.borrow_mut()
.remove(&stream_id)
.ok_or_else(unknown_host_stream)?;
stream.cancel();
Ok(())
}
pub fn deactivate(&self) {
self.bindings.replace(None);
for (_, stream) in std::mem::take(&mut *self.streams.borrow_mut()) {
stream.cancel();
}
}
fn binding(&self, binding_id: u32) -> Result<JsonHostBinding, RuntimeFailure> {
let bindings = self.bindings.borrow();
let bindings = bindings.as_ref().ok_or(RuntimeFailure::AdmissionClosed)?;
bindings
.get(binding_id as usize)
.cloned()
.ok_or(RuntimeFailure::ProtocolViolation {
capability: JSON_HOST_IMPORTS_ABI_V1,
})
}
fn stream(&self, stream_id: u64) -> Result<Rc<dyn JsonHostStreamSession>, RuntimeFailure> {
self.streams
.borrow()
.get(&stream_id)
.cloned()
.ok_or_else(unknown_host_stream)
}
}
fn validate_host_binding(
codec: &Rc<dyn JsonCapabilityCodec>,
request: Option<&PluginDependencyHandle>,
stream: Option<&PluginStreamDependencyHandle>,
) -> Result<(), RuntimeFailure> {
for (capability, version) in request
.map(|handle| (handle.capability_id(), handle.descriptor_version()))
.into_iter()
.chain(stream.map(|handle| (handle.capability_id(), handle.descriptor_version())))
{
if capability != codec.capability_id() || version != codec.descriptor_version() {
return Err(RuntimeFailure::ProtocolViolation {
capability: codec.capability_id(),
});
}
}
Ok(())
}
fn unknown_host_stream() -> RuntimeFailure {
RuntimeFailure::ProtocolViolation {
capability: "lenso.json-host-imports@1",
}
}
impl JsonStreamFrame {
pub fn decode(
encoded: &str,
capability: &'static str,
) -> Result<JsonStreamItem, RuntimeFailure> {
match serde_json::from_str(encoded)
.map_err(|_| RuntimeFailure::ProtocolViolation { capability })?
{
Self::Message(value) => Ok(JsonStreamItem::Message(value)),
Self::PeerHalfClosed => Ok(JsonStreamItem::PeerHalfClosed),
Self::TerminalSuccess => Ok(JsonStreamItem::Terminal(Ok(()))),
Self::TerminalError(value) => Ok(JsonStreamItem::Terminal(Err(value))),
}
}
}
pub trait JsonStreamSessionTransport: std::fmt::Debug + 'static {
fn send(
self: Rc<Self>,
message_json: String,
) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>>;
fn receive(
self: Rc<Self>,
) -> futures::future::LocalBoxFuture<'static, Result<JsonStreamItem, RuntimeFailure>>;
fn close_send(
self: Rc<Self>,
) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>>;
fn cancel(&self);
}
pub type JsonStreamOpenFuture = futures::future::LocalBoxFuture<
'static,
Result<Result<Rc<dyn JsonStreamSessionTransport>, Value>, RuntimeFailure>,
>;
pub trait JsonStreamTransport: std::fmt::Debug + 'static {
fn open(
self: Rc<Self>,
capability: String,
operation: String,
request_json: String,
context: InvocationContext,
) -> JsonStreamOpenFuture;
}
pub fn json_request_endpoints<T: JsonRequestTransport>(
transport: Rc<T>,
codecs: Vec<Rc<dyn JsonCapabilityCodec>>,
) -> Vec<Rc<dyn NativeRequestEndpoint>> {
let transport: Rc<dyn JsonRequestTransport> = transport;
codecs
.into_iter()
.filter(|codec| !codec.request_operations().is_empty())
.map(|codec| {
Rc::new(JsonRequestEndpoint {
transport: transport.clone(),
codec,
}) as Rc<dyn NativeRequestEndpoint>
})
.collect()
}
pub fn json_stream_endpoints<T: JsonStreamTransport>(
transport: Rc<T>,
codecs: Vec<Rc<dyn JsonCapabilityCodec>>,
) -> Vec<Rc<dyn NativeStreamEndpoint>> {
let transport: Rc<dyn JsonStreamTransport> = transport;
codecs
.into_iter()
.filter(|codec| !codec.stream_operations().is_empty())
.map(|codec| {
Rc::new(JsonStreamEndpoint {
transport: transport.clone(),
codec,
}) as Rc<dyn NativeStreamEndpoint>
})
.collect()
}
#[derive(Debug)]
struct JsonStreamEndpoint {
transport: Rc<dyn JsonStreamTransport>,
codec: Rc<dyn JsonCapabilityCodec>,
}
impl NativeStreamEndpoint for JsonStreamEndpoint {
fn capability_id(&self) -> &'static str {
self.codec.capability_id()
}
fn descriptor_version(&self) -> &'static str {
self.codec.descriptor_version()
}
fn operations(&self) -> &'static [&'static str] {
self.codec.stream_operations()
}
fn open(
&self,
operation: &str,
request: Box<dyn Any>,
context: InvocationContext,
) -> futures::future::LocalBoxFuture<
'static,
Result<Result<Box<dyn NativeStreamSession>, Box<dyn Any>>, RuntimeFailure>,
> {
let transport = self.transport.clone();
let codec = self.codec.clone();
let operation = operation.to_owned();
Box::pin(async move {
if !codec.stream_operations().contains(&operation.as_str()) {
return Err(unknown_operation(codec.capability_id(), &operation));
}
let request = codec.encode_stream_open(&operation, request.as_ref())?;
let request_json =
serde_json::to_string(&request).map_err(|_| RuntimeFailure::ProtocolViolation {
capability: codec.capability_id(),
})?;
match transport
.open(
codec.capability_id().to_owned(),
operation.clone(),
request_json,
context,
)
.await?
{
Ok(session) => Ok(Ok(Box::new(JsonStreamSession {
session,
codec,
operation,
}) as Box<dyn NativeStreamSession>)),
Err(error) => codec.decode_stream_domain_error(&operation, error).map(Err),
}
})
}
}
#[derive(Debug)]
struct JsonStreamSession {
session: Rc<dyn JsonStreamSessionTransport>,
codec: Rc<dyn JsonCapabilityCodec>,
operation: String,
}
impl NativeStreamSession for JsonStreamSession {
fn send(
&self,
message: Box<dyn Any>,
) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>> {
let encoded = self
.codec
.encode_stream_message(&self.operation, message.as_ref())
.and_then(|value| {
serde_json::to_string(&value).map_err(|_| RuntimeFailure::ProtocolViolation {
capability: self.codec.capability_id(),
})
});
let session = self.session.clone();
Box::pin(async move { session.send(encoded?).await })
}
fn receive(
&self,
) -> futures::future::LocalBoxFuture<'static, Result<NativeStreamItem, RuntimeFailure>> {
let session = self.session.clone();
let codec = self.codec.clone();
let operation = self.operation.clone();
Box::pin(async move {
match session.receive().await? {
JsonStreamItem::Message(value) => codec
.decode_stream_message(&operation, value)
.map(NativeStreamItem::Message),
JsonStreamItem::PeerHalfClosed => Ok(NativeStreamItem::PeerHalfClosed),
JsonStreamItem::Terminal(Ok(())) => Ok(NativeStreamItem::Terminal(Ok(()))),
JsonStreamItem::Terminal(Err(value)) => codec
.decode_stream_domain_error(&operation, value)
.map(|error| NativeStreamItem::Terminal(Err(error))),
}
})
}
fn close_send(&self) -> futures::future::LocalBoxFuture<'static, Result<(), RuntimeFailure>> {
self.session.clone().close_send()
}
fn cancel(&self) {
self.session.cancel();
}
}
#[derive(Debug)]
struct JsonRequestEndpoint {
transport: Rc<dyn JsonRequestTransport>,
codec: Rc<dyn JsonCapabilityCodec>,
}
impl NativeRequestEndpoint for JsonRequestEndpoint {
fn capability_id(&self) -> &'static str {
self.codec.capability_id()
}
fn descriptor_version(&self) -> &'static str {
self.codec.descriptor_version()
}
fn operations(&self) -> &'static [&'static str] {
self.codec.request_operations()
}
fn invoke(
&self,
operation: &str,
request: Box<dyn Any>,
context: InvocationContext,
) -> futures::future::LocalBoxFuture<
'static,
Result<Result<Box<dyn Any>, Box<dyn Any>>, RuntimeFailure>,
> {
let transport = self.transport.clone();
let codec = self.codec.clone();
let operation = operation.to_owned();
Box::pin(async move {
if !codec.request_operations().contains(&operation.as_str()) {
return Err(RuntimeFailure::UnknownOperation {
capability: codec.capability_id(),
operation,
});
}
let request = codec.encode_request(&operation, request.as_ref())?;
let request =
serde_json::to_string(&request).map_err(|_| RuntimeFailure::ProtocolViolation {
capability: codec.capability_id(),
})?;
match transport
.invoke(
codec.capability_id().to_owned(),
operation.clone(),
request,
context,
)
.await?
{
JsonInvocationOutcome::Success(value) => {
codec.decode_response(&operation, value).map(Ok)
}
JsonInvocationOutcome::DomainError(value) => {
codec.decode_domain_error(&operation, value).map(Err)
}
}
})
}
}
pub fn codecs_for_instance(
instance: &PluginInstancePlan,
codecs: &BTreeMap<String, Rc<dyn JsonCapabilityCodec>>,
) -> Result<Vec<Rc<dyn JsonCapabilityCodec>>, RuntimeFailure> {
let mut selected = Vec::with_capacity(instance.provided_capabilities().len());
for descriptor in instance.provided_capabilities() {
if !descriptor.event_operations().is_empty() {
return Err(RuntimeFailure::InvalidResolvedPlan {
detail: format!(
"Execution class `{}` does not support Event endpoints",
instance.execution_class()
),
});
}
let codec = codecs.get(descriptor.capability_id()).ok_or_else(|| {
RuntimeFailure::InvalidResolvedPlan {
detail: format!(
"no generated codec for Capability `{}`",
descriptor.capability_id()
),
}
})?;
let request_operations: Vec<_> = codec
.request_operations()
.iter()
.map(|operation| (*operation).to_owned())
.collect();
let stream_operations: Vec<_> = codec
.stream_operations()
.iter()
.map(|operation| (*operation).to_owned())
.collect();
let expected_request: Vec<_> = descriptor
.request_operations()
.into_iter()
.map(str::to_owned)
.collect();
let expected_stream: Vec<_> = descriptor
.stream_operations()
.into_iter()
.map(str::to_owned)
.collect();
if codec.descriptor_version() != descriptor.descriptor_version()
|| request_operations != expected_request
|| stream_operations != expected_stream
{
return Err(RuntimeFailure::ProtocolViolation {
capability: codec.capability_id(),
});
}
selected.push(codec.clone());
}
Ok(selected)
}
pub fn codecs_for_requirements(
instance: &PluginInstancePlan,
codecs: &BTreeMap<String, Rc<dyn JsonCapabilityCodec>>,
) -> Result<Vec<Rc<dyn JsonCapabilityCodec>>, RuntimeFailure> {
let mut selected = Vec::with_capacity(instance.required_capabilities().len());
for requirement in instance.required_capabilities() {
let codec = codecs.get(requirement.capability_id()).ok_or_else(|| {
RuntimeFailure::InvalidResolvedPlan {
detail: format!(
"no generated guest import codec for Capability `{}`",
requirement.capability_id()
),
}
})?;
if codec.descriptor_version() != requirement.descriptor_version() {
return Err(RuntimeFailure::ProtocolViolation {
capability: codec.capability_id(),
});
}
selected.push(codec.clone());
}
Ok(selected)
}
pub fn prepare_request_app(
plan: &ResolvedAppPlan,
execution_class: &ExecutionClassId,
generations: BTreeMap<String, PreparedNativePlugin>,
) -> Result<PreparedNativeApp, RuntimeFailure> {
let selected_instances = plan
.plugin_instances()
.iter()
.filter(|instance| instance.execution_class() == execution_class)
.map(|instance| instance.instance_key().to_owned())
.collect::<std::collections::BTreeSet<_>>();
let mut endpoints = BTreeMap::new();
let mut stream_endpoints = BTreeMap::new();
for (instance_key, generation) in &generations {
for endpoint in generation.endpoints() {
let identity = (instance_key.clone(), endpoint.capability_id().to_owned());
if endpoints.insert(identity, endpoint.clone()).is_some() {
return Err(RuntimeFailure::InvalidResolvedPlan {
detail: format!("duplicate request endpoint on Instance `{instance_key}`"),
});
}
}
for endpoint in generation.stream_endpoints() {
let identity = (instance_key.clone(), endpoint.capability_id().to_owned());
if stream_endpoints
.insert(identity, endpoint.clone())
.is_some()
{
return Err(RuntimeFailure::InvalidResolvedPlan {
detail: format!("duplicate stream endpoint on Instance `{instance_key}`"),
});
}
}
}
for instance in plan
.plugin_instances()
.iter()
.filter(|instance| selected_instances.contains(instance.instance_key()))
{
if !generations.contains_key(instance.instance_key()) {
return Err(RuntimeFailure::InvalidResolvedPlan {
detail: format!("Adapter omitted Instance `{}`", instance.instance_key()),
});
}
}
let mut bindings = Vec::new();
let mut stream_bindings = Vec::new();
for binding in plan.capability_bindings() {
let key = (
binding.provider_instance().to_owned(),
binding.capability_id().to_owned(),
);
let request_endpoint = endpoints.get(&key);
let stream_endpoint = stream_endpoints.get(&key);
if let Some(endpoint) = request_endpoint {
bindings.push(PreparedBinding::new(
binding.consumer_instance(),
binding.provider_instance(),
endpoint.clone(),
));
}
if let Some(endpoint) = stream_endpoint {
stream_bindings.push(PreparedStreamBinding::new(
binding.consumer_instance(),
binding.provider_instance(),
endpoint.clone(),
));
}
if request_endpoint.is_none()
&& stream_endpoint.is_none()
&& selected_instances.contains(binding.provider_instance())
{
return Err(RuntimeFailure::InvalidResolvedPlan {
detail: format!(
"Adapter omitted Capability `{}` endpoint for Instance `{}`",
binding.capability_id(),
binding.provider_instance()
),
});
}
}
Ok(PreparedNativeApp::new(bindings, generations).with_stream_bindings(stream_bindings))
}
pub fn require_operation(
codecs: &BTreeMap<String, Rc<dyn JsonCapabilityCodec>>,
capability_id: &str,
operation: &str,
) -> Result<Rc<dyn JsonCapabilityCodec>, RuntimeFailure> {
let codec =
codecs
.get(capability_id)
.cloned()
.ok_or_else(|| RuntimeFailure::InvalidResolvedPlan {
detail: format!("no generated codec for Capability `{capability_id}`"),
})?;
if !codec.request_operations().contains(&operation) {
return Err(RuntimeFailure::UnknownOperation {
capability: codec.capability_id(),
operation: operation.to_owned(),
});
}
Ok(codec)
}
fn unknown_operation(capability: &'static str, operation: &str) -> RuntimeFailure {
RuntimeFailure::UnknownOperation {
capability,
operation: operation.to_owned(),
}
}
fn validate_digest(digest: &str) -> Result<(), RuntimeFailure> {
let valid = digest.strip_prefix("sha256:").is_some_and(|hex| {
hex.len() == 64
&& hex
.bytes()
.all(|byte| byte.is_ascii_hexdigit() && !byte.is_ascii_uppercase())
});
if valid {
Ok(())
} else {
Err(RuntimeFailure::InvalidResolvedPlan {
detail: format!("invalid canonical SHA-256 digest `{digest}`"),
})
}
}
fn invalid_artifact(path: &Path, error: impl std::fmt::Display) -> RuntimeFailure {
RuntimeFailure::InvalidResolvedPlan {
detail: format!("cannot read Artifact `{}`: {error}", path.display()),
}
}
#[cfg(test)]
mod tests {
use std::io::Write;
use super::*;
#[test]
fn artifact_handle_rejects_digest_drift() {
let mut file = tempfile::NamedTempFile::new().unwrap();
file.write_all(b"first").unwrap();
let digest = format!("sha256:{}", hex::encode(Sha256::digest(b"first")));
let handle = ArtifactHandle::open(file.path(), &digest, 5).unwrap();
file.as_file_mut().set_len(0).unwrap();
file.write_all(b"other").unwrap();
assert!(handle.read_verified().is_err());
}
}