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