Skip to main content

lenso_runtime_codec/
lib.rs

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