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