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