Skip to main content

smolvm_protocol/
lib.rs

1//! Protocol types for smolvm host-guest communication.
2//!
3//! This crate defines the wire protocol for vsock communication between
4//! the smolvm host and the guest agent (smolvm-agent).
5//!
6//! # Protocol Overview
7//!
8//! Communication uses JSON-encoded messages over vsock. Each message is
9//! prefixed with a 4-byte big-endian length header.
10//!
11//! ```text
12//! +----------------+-------------------+
13//! | Length (4 BE)  | JSON payload      |
14//! +----------------+-------------------+
15//! ```
16
17#![deny(missing_docs)]
18
19use serde::{Deserialize, Serialize};
20
21/// One S3-compatible bucket to mount inside the workload container.
22///
23/// Structured rather than a shell command: the agent mounts it natively, so
24/// nothing has to be installed in the image and no command is interpolated.
25#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
26pub struct S3Volume {
27    /// Service endpoint, e.g. `https://s3.us-east-1.amazonaws.com` or a
28    /// self-hosted `http://minio:9000`.
29    pub endpoint: String,
30    /// Region used for request signing.
31    pub region: String,
32    /// Bucket to mount.
33    pub bucket: String,
34    /// Optional key prefix, so a mount can expose one sub-tree of a bucket.
35    #[serde(default)]
36    pub prefix: String,
37    /// Absolute path inside the container where the bucket appears.
38    pub mountpoint: String,
39    /// Mount read-only; the kernel then rejects writes before they reach us.
40    #[serde(default)]
41    pub read_only: bool,
42    /// Access key. Absent (with the secret) means anonymous access.
43    #[serde(default, skip_serializing_if = "Option::is_none")]
44    pub access_key_id: Option<String>,
45    /// Secret key paired with `access_key_id`.
46    #[serde(default, skip_serializing_if = "Option::is_none")]
47    pub secret_access_key: Option<String>,
48    /// Session token, when temporary credentials are in use.
49    #[serde(default, skip_serializing_if = "Option::is_none")]
50    pub session_token: Option<String>,
51}
52
53pub mod credentials;
54pub mod forkpoint;
55pub mod guest_env;
56pub mod image_ref;
57pub mod intercept;
58pub mod publish_socket;
59pub mod retry;
60pub mod secrets;
61
62pub use credentials::{CredentialBinding, CredentialPolicy};
63pub use image_ref::{image_repo, normalize_image_ref};
64pub use intercept::InterceptEndpoint;
65pub use secrets::{SecretRef, SecretSourceKind};
66
67/// Serde helper for encoding `Vec<u8>` as a base64 string in JSON.
68///
69/// Without this, serde_json serializes `Vec<u8>` as a JSON array of numbers
70/// (e.g., `[104,101,108,108,111]`), which inflates binary data by ~4x.
71/// Base64 encoding reduces this to ~1.33x.
72pub mod base64_bytes {
73    use base64::{engine::general_purpose::STANDARD, Engine};
74    use serde::{Deserialize, Deserializer, Serializer};
75
76    /// Serialize `Vec<u8>` as a base64 string.
77    pub fn serialize<S: Serializer>(data: &[u8], serializer: S) -> Result<S::Ok, S::Error> {
78        serializer.serialize_str(&STANDARD.encode(data))
79    }
80
81    /// Deserialize a base64 string into `Vec<u8>`.
82    pub fn deserialize<'de, D: Deserializer<'de>>(deserializer: D) -> Result<Vec<u8>, D::Error> {
83        let s = String::deserialize(deserializer)?;
84        STANDARD.decode(&s).map_err(serde::de::Error::custom)
85    }
86}
87
88/// Protocol version.
89pub const PROTOCOL_VERSION: u32 = 1;
90
91/// The agent can freeze internal filesystems before acknowledging shutdown.
92pub const QUIESCED_SHUTDOWN_CAPABILITY: &str = "quiesced-shutdown-v1";
93
94/// virtiofs tag under which the host exposes the Rosetta 2 Linux runtime to the
95/// guest. Shared host↔guest so the launcher's `krun_add_virtiofs` tag and the
96/// guest agent's `mount -t virtiofs` source can't drift apart.
97pub const ROSETTA_TAG: &str = "rosetta";
98
99/// Guest mount point for the Rosetta 2 Linux runtime. The ptrace wrapper execs
100/// `<ROSETTA_GUEST_PATH>/rosetta` (the translator), so this path is baked into
101/// both the wrapper and the `binfmt_misc` registration.
102pub const ROSETTA_GUEST_PATH: &str = "/mnt/rosetta";
103
104/// Maximum frame size (32 MB - layer exports use chunked streaming).
105pub const MAX_FRAME_SIZE: u32 = 32 * 1024 * 1024;
106
107/// Chunk size for streaming layer data (~16 MB raw, ~21 MB as base64 JSON).
108pub const LAYER_CHUNK_SIZE: usize = 16 * 1024 * 1024;
109
110/// Files at or below this size are written with a single `FileWrite`
111/// message. Larger files must stream via
112/// `FileWriteBegin` + `FileWriteChunk` so no single frame approaches
113/// [`MAX_FRAME_SIZE`] (base64 + JSON inflation is ~1.4x).
114///
115/// Chosen to keep the single-shot frame comfortably under the frame
116/// limit while preserving the fast-path latency for small config
117/// files / scripts / keys.
118pub const FILE_WRITE_SINGLE_SHOT_MAX: usize = 1024 * 1024;
119
120/// Payload bytes per streaming upload chunk. Deliberately small —
121/// equal to [`FILE_WRITE_SINGLE_SHOT_MAX`] — so each chunk's encoded
122/// frame (~1.4 MB) fits inside typical kernel Unix-socket send
123/// buffers (`SO_SNDBUF` defaults on the order of 200–256 KiB but
124/// can grow). Larger chunks would force `write_all` to spin waiting
125/// for the agent to drain, and any latency spike trips the 10 s
126/// write timeout with `EAGAIN` — exactly the failure David
127/// reproduced before this fix landed.
128///
129/// Note: [`LAYER_CHUNK_SIZE`] is 16 MiB for agent→host (download)
130/// streaming, which works because the host side of the socket has
131/// more headroom than the guest side. Upload streaming is the
132/// asymmetric case and needs a smaller chunk.
133pub const FILE_WRITE_CHUNK_SIZE: usize = FILE_WRITE_SINGLE_SHOT_MAX;
134
135// The single-shot threshold must be <= the chunk size. They can be equal (a
136// 1 MiB file is a single shot; a 1 MiB + 1 byte file streams as two chunks),
137// but SINGLE_SHOT > CHUNK would be incoherent — a file slightly over the shot
138// threshold would have to stream as a single oversized chunk. Enforced at
139// compile time so a bad edit fails the build rather than a test.
140const _: () = assert!(FILE_WRITE_SINGLE_SHOT_MAX <= FILE_WRITE_CHUNK_SIZE);
141
142/// Hard ceiling on a single file transfer in either direction.
143///
144/// On the write path: enforced at `FileWriteBegin` by the agent —
145/// `total_size > FILE_TRANSFER_MAX_TOTAL` is rejected before any
146/// staging file is created.
147///
148/// On the read path: enforced by the host's `read_file` loop —
149/// after the first chunk that pushes the accumulated total past the
150/// cap, the call bails with an error and the partial buffer is
151/// dropped. This protects the host process from OOM if the guest
152/// (compromised or merely buggy) streams unbounded data.
153///
154/// 4 GiB matches the order-of-magnitude of the default overlay disk and the
155/// `gpu_vram_mib` cap. It is the compiled-in default. The host-side transfer
156/// paths (read/export, including `pack create --from-vm`) raise it at runtime
157/// via `SMOLVM_FILE_TRANSFER_MAX_BYTES` (see
158/// `agent::client::file_transfer_max_total`) so a VM snapshot whose overlay
159/// carries a large dependency tree (e.g. a torch + CUDA-wheels environment is
160/// ~5 GiB) can be packed without lowering the DoS bound for everyone else. The
161/// guest agent's write-path check runs inside the VM and cannot see that host
162/// env var, so it always enforces this const. Callers that need to move larger
163/// blobs routinely should stage via a virtiofs mount instead of `cp`.
164pub const FILE_TRANSFER_MAX_TOTAL: u64 = 4 * 1024 * 1024 * 1024;
165
166/// Filename of the virtiofs-visible marker the agent creates when it is
167/// ready to accept vsock connections.
168///
169/// The host polls for this file through its virtiofs mount of the guest
170/// rootfs. The agent writes it (and optionally a symlink from `/oldroot/`)
171/// during deferred init, just before opening the vsock listener.
172///
173/// Both sides must agree on this name; keeping it here prevents silent drift.
174pub const AGENT_READY_MARKER: &str = ".smolvm-ready";
175
176/// Well-known vsock ports.
177pub mod ports {
178    /// Control channel for workload VMs.
179    pub const WORKLOAD_CONTROL: u32 = 5000;
180    /// Log streaming from workload VMs.
181    pub const WORKLOAD_LOGS: u32 = 5001;
182    /// Agent control port (for OCI operations and management).
183    pub const AGENT_CONTROL: u32 = 6000;
184    /// SSH agent forwarding (host SSH_AUTH_SOCK bridged to guest).
185    pub const SSH_AGENT: u32 = 6001;
186    /// DNS filtering proxy (guest forwards DNS queries to host for filtering).
187    pub const DNS_FILTER: u32 = 6002;
188    /// Docker socket bridge: the guest listens on this vsock port and proxies
189    /// each connection to the in-guest `/var/run/docker.sock`, so the host can
190    /// reach the guest's Docker daemon over a host-side Unix socket
191    /// (`DOCKER_HOST=unix://…`). Inbound (host connects in), like the agent
192    /// control channel — unlike the outbound SSH/DNS/CUDA bridges.
193    pub const DOCKER: u32 = 6003;
194    /// Readiness doorbell: the guest connects OUT to this host port the instant it
195    /// finishes init, and the host's `accept()` fires — an event-driven readiness
196    /// signal that needs no writable filesystem, DAX window, or extra virtio
197    /// device. Outbound (guest connects in), like the SSH/DNS/CUDA bridges. The
198    /// primary readiness signal; the ready-marker file and the control-channel
199    /// ping remain as fallbacks. A fork clone resumes past the ring, so it never
200    /// dials this — clones detect readiness by pinging the restored agent.
201    pub const AGENT_READY: u32 = 6004;
202    /// CUDA-over-vsock (experimental): guest CUDA client forwards Driver-API
203    /// calls to a host CUDA server that runs them on the host NVIDIA GPU.
204    pub const CUDA: u32 = 7000;
205
206    /// Base vsock port for user-published Unix-socket bridges
207    /// (`--expose-socket` / `--mount-socket`). Each published socket is assigned
208    /// `PUBLISH_SOCKET_BASE + index`. Kept clear of the fixed ports above (and of
209    /// CUDA at 7000) so a reasonable number of sockets never collides.
210    pub const PUBLISH_SOCKET_BASE: u32 = 6100;
211
212    /// Maximum number of user-published sockets per VM. Bounds the vsock-port
213    /// window (`6100..6100+MAX`) below CUDA's 7000.
214    pub const PUBLISH_SOCKET_MAX: usize = 64;
215}
216
217/// vsock CID constants.
218pub mod cid {
219    /// Host CID (always 2).
220    pub const HOST: u32 = 2;
221    /// Guest CID (always 3 for the first/only guest).
222    pub const GUEST: u32 = 3;
223    /// Any CID (for listening).
224    pub const ANY: u32 = u32::MAX;
225}
226
227/// fsnotify event masks, mirroring the kernel's `FS_*` bits in
228/// `include/linux/fsnotify_backend.h`. Shared by the host watcher (which maps a
229/// host filesystem event to one of these) and the guest agent (which forwards
230/// the raw bits to `/proc/smolvm-fsnotify`). Only the subset relevant to
231/// file-watching tools is defined.
232pub mod fsnotify_mask {
233    /// File was modified.
234    pub const FS_MODIFY: u32 = 0x0000_0002;
235    /// Metadata changed (chmod/chown/utimes).
236    pub const FS_ATTRIB: u32 = 0x0000_0004;
237    /// Writable file was closed.
238    pub const FS_CLOSE_WRITE: u32 = 0x0000_0008;
239    /// File was moved away from the watched dir.
240    pub const FS_MOVED_FROM: u32 = 0x0000_0040;
241    /// File was moved into the watched dir.
242    pub const FS_MOVED_TO: u32 = 0x0000_0080;
243    /// Subfile was created.
244    pub const FS_CREATE: u32 = 0x0000_0100;
245    /// Subfile was deleted.
246    pub const FS_DELETE: u32 = 0x0000_0200;
247    /// Event occurred against a directory.
248    pub const FS_ISDIR: u32 = 0x4000_0000;
249}
250
251/// A single host-originated filesystem change to replay into the guest.
252#[derive(Debug, Clone, Serialize, Deserialize)]
253pub struct FsNotifyEvent {
254    /// Guest-side absolute path the event occurred on (virtiofs staging path).
255    pub path: String,
256    /// `fsnotify_mask::FS_*` bitmask for the event.
257    pub mask: u32,
258}
259
260// ============================================================================
261// Agent Protocol (OCI Operations)
262// ============================================================================
263
264fn shutdown_progress_disabled(progress: &bool) -> bool {
265    !progress
266}
267
268#[cfg(test)]
269mod shutdown_compat_tests {
270    use super::*;
271
272    #[test]
273    fn legacy_shutdown_wire_shape_is_unchanged() {
274        let wire = r#"{"method":"shutdown"}"#;
275        assert!(matches!(
276            serde_json::from_str::<AgentRequest>(wire).unwrap(),
277            AgentRequest::Shutdown { progress: false }
278        ));
279        assert_eq!(
280            serde_json::to_string(&AgentRequest::Shutdown { progress: false }).unwrap(),
281            wire
282        );
283    }
284
285    #[test]
286    fn older_agent_can_accept_opt_in_request() {
287        #[derive(Deserialize)]
288        #[serde(tag = "method", rename_all = "snake_case")]
289        enum LegacyRequest {
290            Shutdown,
291        }
292        let wire = serde_json::to_string(&AgentRequest::Shutdown { progress: true }).unwrap();
293        assert!(matches!(
294            serde_json::from_str::<LegacyRequest>(&wire).unwrap(),
295            LegacyRequest::Shutdown
296        ));
297    }
298
299    #[test]
300    fn traced_shutdown_preserves_progress_opt_in() {
301        let envelope = Envelope::with_trace_id(
302            AgentRequest::Shutdown { progress: true },
303            Some("shutdown-test".into()),
304        );
305        let decoded: Envelope<AgentRequest> =
306            decode_message(&encode_message(&envelope).unwrap()).unwrap();
307        assert!(matches!(
308            decoded.body,
309            AgentRequest::Shutdown { progress: true }
310        ));
311        assert_eq!(decoded.trace_id.as_deref(), Some("shutdown-test"));
312    }
313}
314
315/// Agent request types (for image management and OCI operations).
316#[derive(Debug, Clone, Serialize, Deserialize)]
317#[serde(tag = "method", rename_all = "snake_case")]
318pub enum AgentRequest {
319    /// Ping to check if agent is alive.
320    Ping,
321
322    /// Inject host-originated fsnotify events into the guest.
323    ///
324    /// virtiofs does not deliver host-side file changes to the guest as
325    /// fsnotify/inotify events, so inotify-based hot-reload (Vite, webpack,
326    /// nodemon) never fires when a mounted file is edited on the host. The host
327    /// watches the mount source and sends the resulting events here; the agent
328    /// writes them to `/proc/smolvm-fsnotify`, which fires the matching event on
329    /// the guest inode so watchers on the (bind-mounted) container path wake up.
330    /// Each `path` is a guest-side absolute path (the virtiofs staging path),
331    /// `mask` an `fsnotify_mask::FS_*` bitmask.
332    FsNotify {
333        /// Host-originated filesystem changes to replay as guest fsnotify events.
334        #[serde(default)]
335        events: Vec<FsNotifyEvent>,
336    },
337
338    /// Pull an OCI image and extract layers.
339    Pull {
340        /// Image reference (e.g., "alpine:latest", "docker.io/library/ubuntu:22.04").
341        image: String,
342        /// OCI platform to pull (e.g., "linux/arm64", "linux/amd64").
343        oci_platform: Option<String>,
344        /// Optional registry authentication credentials.
345        #[serde(default, skip_serializing_if = "Option::is_none")]
346        auth: Option<RegistryAuth>,
347        /// Proxy URL applied to the registry client (sets HTTP_PROXY and HTTPS_PROXY).
348        #[serde(default, skip_serializing_if = "Option::is_none")]
349        proxy: Option<String>,
350        /// Comma-separated NO_PROXY list of hosts/CIDRs that bypass the proxy.
351        #[serde(default, skip_serializing_if = "Option::is_none")]
352        no_proxy: Option<String>,
353    },
354
355    /// Query if an image exists locally.
356    Query {
357        /// Image reference.
358        image: String,
359    },
360
361    /// List all cached images.
362    ListImages,
363
364    /// Run garbage collection on unused layers.
365    GarbageCollect {
366        /// If true, only report what would be deleted.
367        dry_run: bool,
368        /// If true, delete all image manifests and configs first,
369        /// making all layers unreferenced so they get collected.
370        #[serde(default)]
371        purge_all: bool,
372    },
373
374    /// Prepare overlay rootfs for a workload.
375    PrepareOverlay {
376        /// Image reference.
377        image: String,
378        /// Unique workload ID for the overlay.
379        workload_id: String,
380    },
381
382    /// Clean up overlay rootfs for a workload.
383    CleanupOverlay {
384        /// Workload ID to clean up.
385        workload_id: String,
386    },
387
388    /// Format the storage disk (first-time setup).
389    FormatStorage,
390
391    /// Get storage disk status.
392    StorageStatus,
393
394    /// Report the guest's own view of machine memory, read from
395    /// `/proc/meminfo`.
396    ///
397    /// The host cannot answer this. macOS charges the VMM's `phys_footprint`
398    /// as `internal + compressed`, and the compressed part is counted at the
399    /// pages' *uncompressed* size, so an idle guest whose memory has been
400    /// compressed reports a footprint several times the bytes it actually
401    /// occupies. Only the guest's allocator knows what the machine is using.
402    MemoryStatus,
403
404    /// Test network connectivity directly from the agent (not via chroot).
405    /// Used to debug TSI networking.
406    NetworkTest {
407        /// URL to test (e.g., "http://1.1.1.1")
408        url: String,
409    },
410
411    /// Shutdown the agent.
412    Shutdown {
413        /// Opt into progress frames while storage is being synchronized.
414        #[serde(default, skip_serializing_if = "shutdown_progress_disabled")]
415        progress: bool,
416    },
417
418    /// Export a layer as a tar archive.
419    ///
420    /// Used by `smolvm pack` to extract OCI layers for packaging.
421    /// The agent streams the layer tar data back via LayerData responses.
422    ExportLayer {
423        /// Image digest (sha256:...).
424        image_digest: String,
425        /// Layer index (0-based).
426        layer_index: usize,
427    },
428
429    /// Merge a stack of directories into a single tar archive.
430    ///
431    /// `pack create --from-vm` ships the machine's image layers plus whatever the
432    /// container has written since as ONE flattened layer. The merge has to apply
433    /// whiteouts and opaque markers exactly as the runtime would, so it is done
434    /// with a read-only overlay mount rather than a file copy.
435    ///
436    /// The agent owns this rather than the host driving `mount(8)` over `VmExec`,
437    /// because `mount(8)` rejects a `lowerdir=` value beyond ~255 bytes — about
438    /// three OCI layer paths — while the agent can append each layer separately
439    /// through the same `fsconfig` path the runtime container mount uses.
440    FlattenLayers {
441        /// Directories to merge, bottom -> top. Entries that are missing or empty
442        /// are dropped, so a caller may include a container overlay's upper dir
443        /// without first checking whether the machine ever wrote to it.
444        lowerdirs: Vec<String>,
445        /// Guest path to write the tar archive to, or `None` to stream the
446        /// archive straight back as `DataChunk` responses.
447        ///
448        /// Streaming is what a large export wants. The flattened tar is as large
449        /// as the image it came from, so staging it in the guest means the disk
450        /// has to hold both the expanded rootfs and a second full copy of it as
451        /// an archive — the export sizes that disk by guessing, and a big enough
452        /// image runs it out of space.
453        #[serde(default)]
454        output: Option<String>,
455    },
456
457    /// Wait until the workload has declared a branchpoint, returning the ready
458    /// marker's contents (its profile lines) in `data.contents`.
459    BranchpointWait {
460        /// Give up after this many milliseconds.
461        timeout_ms: u64,
462    },
463    /// Put a negotiated helper into its restore-safe loop before capture.
464    BranchpointArm,
465    /// Return a parked source to its ordinary wait after capture.
466    BranchpointPark,
467    /// Release a restored clone: write the release marker for the generation
468    /// recorded in its ready marker, carrying the clone's identity, and wait
469    /// for the helper to acknowledge.
470    BranchpointRelease {
471        /// The clone's parameters in dotenv form, appended to the marker.
472        #[serde(default)]
473        env_dotenv: Option<String>,
474    },
475    /// Assign and release a held clone in one idempotent step: claim it with
476    /// `activation_token` (a retry with the same token completes a partial
477    /// commit; a different token is refused), install the per-clone
478    /// parameters, then write the release marker.
479    BranchpointActivate {
480        /// Per-clone parameters, dotenv form, written to `env_path`.
481        env_dotenv: String,
482        /// The same parameters in shell-sourceable form, written to `branch_env_path`.
483        env_sourceable: String,
484        /// Guest path of the dotenv file (under the clone's overlay for image machines).
485        env_path: String,
486        /// Guest path of the sourceable file.
487        branch_env_path: String,
488        /// Directory that must already exist, typically the clone's merged
489        /// overlay root; activation refuses rather than fabricating it.
490        require_dir: Option<String>,
491        /// Directory to create before writing the env files.
492        env_dir: String,
493        /// Token identifying this activation attempt.
494        activation_token: String,
495    },
496    /// Wait until a released clone's workload publishes its worker-ready token.
497    BranchpointWaitWorkerReady {
498        /// The token the workload must publish; any other is a mismatch.
499        token: String,
500        /// Give up after this many milliseconds.
501        timeout_ms: u64,
502    },
503
504    /// Execute a command directly in the VM (not in a container).
505    ///
506    /// This runs the command in the agent's Alpine rootfs without any
507    /// container isolation. Useful for VM-level operations and debugging.
508    VmExec {
509        /// Command and arguments.
510        command: Vec<String>,
511        /// Environment variables.
512        #[serde(default)]
513        env: Vec<(String, String)>,
514        /// Working directory in the VM.
515        workdir: Option<String>,
516        /// Timeout in milliseconds.
517        #[serde(default)]
518        timeout_ms: Option<u64>,
519        /// Interactive mode - stream I/O instead of buffering.
520        #[serde(default)]
521        interactive: bool,
522        /// Allocate a pseudo-TTY for the command.
523        #[serde(default)]
524        tty: bool,
525        /// Background mode - spawn and return PID immediately without waiting.
526        #[serde(default)]
527        background: bool,
528        /// Data to pipe to the command's stdin.
529        #[serde(default)]
530        stdin_data: Option<String>,
531    },
532
533    /// Run a command in an image's rootfs.
534    ///
535    /// This prepares an overlay, chroots into it, and executes the command.
536    /// Returns stdout, stderr, and exit code when the command completes.
537    Run {
538        /// Image reference (must be pulled first).
539        image: String,
540        /// Command and arguments.
541        command: Vec<String>,
542        /// Environment variables.
543        #[serde(default)]
544        env: Vec<(String, String)>,
545        /// Working directory inside the rootfs.
546        workdir: Option<String>,
547        /// User inside the rootfs. If omitted, the OCI image default applies.
548        #[serde(default, skip_serializing_if = "Option::is_none")]
549        user: Option<String>,
550        /// Volume mounts to bind into the container.
551        /// Each tuple is (virtiofs_tag, container_path, read_only).
552        #[serde(default)]
553        mounts: Vec<(String, String, bool)>,
554        /// Timeout in milliseconds. If the command exceeds this duration,
555        /// it will be killed and return exit code 124.
556        #[serde(default)]
557        timeout_ms: Option<u64>,
558        /// Interactive mode - stream I/O instead of buffering.
559        /// When true, output is streamed via Stdout/Stderr responses,
560        /// and stdin can be sent via the Stdin request.
561        #[serde(default)]
562        interactive: bool,
563        /// Allocate a pseudo-TTY for the command.
564        /// Enables terminal features like colors, line editing, and signal handling.
565        #[serde(default)]
566        tty: bool,
567        /// Detached mode — start the container and return immediately with the
568        /// container ID. Only meaningful when `persistent_overlay_id` is set.
569        /// Returns a `Completed` response with `stdout` containing the container ID.
570        #[serde(default)]
571        detached: bool,
572        /// Run the workload as an unprivileged container: restricted capabilities,
573        /// read-only cgroup, and no extra tmpfs. The default (false) is "VM-grade"
574        /// — since the microVM is the isolation boundary, the workload gets a full
575        /// capability set and the mounts an init system needs (so any image, incl.
576        /// systemd, boots). Opt in for defense-in-depth when running untrusted code.
577        #[serde(default)]
578        unprivileged: bool,
579        /// If set, use a persistent overlay that survives across exec sessions.
580        /// The overlay is identified by this ID (typically the machine name)
581        /// and reused on subsequent runs. If not set, an ephemeral overlay is
582        /// created and destroyed after the run.
583        #[serde(default, skip_serializing_if = "Option::is_none")]
584        persistent_overlay_id: Option<String>,
585        /// Data to pipe to the command's stdin (non-interactive runs only).
586        /// The pipe is closed after writing, so the command sees EOF.
587        #[serde(default, skip_serializing_if = "Option::is_none")]
588        stdin_data: Option<String>,
589        /// Spawn the container and return immediately with the crun PID.
590        /// The container runs detached; stdout/stderr go to /dev/null.
591        /// Incompatible with `interactive` and `tty`.
592        #[serde(default)]
593        background: bool,
594        /// S3 volumes to mount into the workload container. The agent mounts
595        /// them between `crun create` and `crun start`, so the workload's first
596        /// instruction already sees them and its command is never rewritten.
597        #[serde(default, skip_serializing_if = "Vec::is_empty")]
598        s3_volumes: Vec<S3Volume>,
599    },
600
601    /// Send stdin data to a running interactive command.
602    Stdin {
603        /// Input data to send to the command's stdin.
604        #[serde(with = "base64_bytes")]
605        data: Vec<u8>,
606    },
607
608    /// Resize the PTY window (for TTY mode).
609    Resize {
610        /// New width in columns.
611        cols: u16,
612        /// New height in rows.
613        rows: u16,
614    },
615
616    // ========================================================================
617    // File I/O
618    // ========================================================================
619    /// Write a file inside the VM in a single message.
620    ///
621    /// Use only for files up to [`FILE_WRITE_SINGLE_SHOT_MAX`]. Larger
622    /// files must stream via [`Self::FileWriteBegin`] +
623    /// [`Self::FileWriteChunk`] to avoid exceeding [`MAX_FRAME_SIZE`]
624    /// after base64 + JSON inflation.
625    FileWrite {
626        /// Absolute path in the VM filesystem.
627        path: String,
628        /// File contents.
629        #[serde(with = "base64_bytes")]
630        data: Vec<u8>,
631        /// File mode (e.g., 0o644). None = default (0644).
632        #[serde(default)]
633        mode: Option<u32>,
634        /// Owner uid to apply after the write. None = leave as written (root).
635        #[serde(default)]
636        uid: Option<u32>,
637        /// Owner gid to apply after the write. None = leave as written (root).
638        #[serde(default)]
639        gid: Option<u32>,
640    },
641
642    /// Open a streaming file upload session on this connection.
643    ///
644    /// Must be followed by one or more [`Self::FileWriteChunk`]
645    /// requests. The final chunk sets `done: true` to finalize.
646    /// Dropping the connection (or sending any non-chunk request)
647    /// before `done` aborts the session and leaves no partial file
648    /// at `path`.
649    ///
650    /// Sessions are per-connection — one session at a time.
651    FileWriteBegin {
652        /// Absolute path in the VM filesystem.
653        path: String,
654        /// File mode (e.g., 0o644). None = default (0644).
655        #[serde(default)]
656        mode: Option<u32>,
657        /// Owner uid to apply on finalize. None = leave as written (root).
658        #[serde(default)]
659        uid: Option<u32>,
660        /// Owner gid to apply on finalize. None = leave as written (root).
661        #[serde(default)]
662        gid: Option<u32>,
663        /// Expected total size in bytes. Rejected if it exceeds
664        /// [`FILE_TRANSFER_MAX_TOTAL`]. The agent uses this for an
665        /// early-fail check only; the actual size written is the sum
666        /// of chunk byte lengths.
667        total_size: u64,
668    },
669
670    /// Append a chunk to the currently open streaming upload.
671    /// If `done` is true, the agent fsyncs and atomically renames the
672    /// staging file onto the target path.
673    FileWriteChunk {
674        /// Chunk bytes. Typically [`FILE_WRITE_CHUNK_SIZE`] except
675        /// for the last chunk.
676        #[serde(with = "base64_bytes")]
677        data: Vec<u8>,
678        /// True on the final chunk; closes and renames the staging
679        /// file. False on intermediate chunks.
680        done: bool,
681    },
682
683    /// Read a file from the VM.
684    FileRead {
685        /// Absolute path in the VM filesystem.
686        path: String,
687    },
688
689    /// List the entries of a guest directory.
690    ///
691    /// Kept separate from [`AgentRequest::FileRead`], which streams bytes: a
692    /// listing is small and structured, and a caller that has to discover what
693    /// exists would otherwise guess names and read a 404 for each miss.
694    ListDirectory {
695        /// Absolute directory path in the VM filesystem.
696        path: String,
697    },
698
699    /// Stream a tar archive of a guest directory without creating a guest-side
700    /// temporary file. Used by staged mounts to batch many small files across
701    /// vsock instead of paying one virtiofs round trip per file.
702    ArchiveDirectory {
703        /// Absolute directory path in the VM filesystem.
704        path: String,
705    },
706
707    /// Create (without starting) a Kubernetes pod container whose rootfs is a
708    /// virtiofs-shared host directory (containerd snapshotter output) and whose
709    /// process definition comes from the host's OCI config. The agent builds
710    /// its crun bundle around the shared rootfs; nothing runs until
711    /// `PodStart`. Part of the containerd shim v2 datapath
712    /// (docs/kubernetes-runtime.md).
713    PodCreate {
714        /// Container ID (containerd task id).
715        id: String,
716        /// Rootfs path relative to the sandbox's shared virtiofs mount. The
717        /// shim boots the sandbox VM with ONE shared dir and bind-mounts each
718        /// container's rootfs under it (virtiofs shares are fixed at boot, but
719        /// pod containers are created afterwards), so the guest resolves this
720        /// as `<sandbox-share-mount>/<rootfs_rel>`.
721        rootfs_rel: String,
722        /// The host OCI runtime spec (config.json bytes). The agent extracts
723        /// process/env/cwd/user/mounts/resources and grafts them onto its own
724        /// guest bundle template; host-specific namespaces/paths are ignored.
725        spec_json: String,
726        /// Allocate a PTY for the init process.
727        #[serde(default)]
728        tty: bool,
729    },
730
731    /// Start a pod container created by `PodCreate` (or an exec process
732    /// registered by `PodExec`), streaming its I/O on THIS connection:
733    /// `Started` → `Stdout`/`Stderr`... → `Exited`. Stdin arrives via `Stdin`
734    /// requests; PTY resize via `Resize`.
735    PodStart {
736        /// Container ID.
737        id: String,
738        /// Exec process to start instead of the init process.
739        #[serde(default, skip_serializing_if = "Option::is_none")]
740        exec_id: Option<String>,
741    },
742
743    /// Register an exec process for a running pod container. Started later by
744    /// `PodStart { exec_id }`.
745    PodExec {
746        /// Container ID.
747        id: String,
748        /// Exec process ID (unique within the container).
749        exec_id: String,
750        /// OCI Process JSON (containerd's ExecProcessRequest spec).
751        process_json: String,
752        /// Allocate a PTY for the exec process.
753        #[serde(default)]
754        tty: bool,
755    },
756
757    /// Signal a pod container's init process (or one exec process).
758    PodSignal {
759        /// Container ID.
760        id: String,
761        /// Exec process to signal instead of init.
762        #[serde(default, skip_serializing_if = "Option::is_none")]
763        exec_id: Option<String>,
764        /// Signal number (SIGKILL = 9, SIGTERM = 15, ...).
765        signal: u32,
766        /// Signal the whole container process group.
767        #[serde(default)]
768        all: bool,
769    },
770
771    /// List PIDs inside a pod container (guest view).
772    PodPids {
773        /// Container ID.
774        id: String,
775    },
776
777    /// Sample a pod container's resource usage (guest view). The agent reads the
778    /// container's process tree from /proc (there is no per-container cgroup); the
779    /// shim maps the reply into containerd's cgroups metrics for CRI stats.
780    PodStats {
781        /// Container ID.
782        id: String,
783    },
784
785    /// Remove a pod container's (or exec process's) guest resources after
786    /// exit: bundle, cgroup, PTY. Exit status was already streamed by
787    /// `PodStart`'s `Exited`.
788    PodDelete {
789        /// Container ID.
790        id: String,
791        /// Exec process to remove instead of the whole container.
792        #[serde(default, skip_serializing_if = "Option::is_none")]
793        exec_id: Option<String>,
794    },
795}
796
797impl AgentRequest {
798    /// A log-safe one-line summary of the request.
799    ///
800    /// This string is written to the machine's console log, which is exposed
801    /// over the logs API — so it must NEVER include credential- or data-bearing
802    /// fields: registry `auth`, `env` (which can carry host-resolved secrets),
803    /// `proxy` (may embed credentials), or `data` (file/stdin bytes). Only the
804    /// variant name plus a non-secret identifier (image) is emitted.
805    ///
806    /// The match is exhaustive with no catch-all on purpose: adding a new
807    /// variant forces a compile error here, so redaction is a deliberate
808    /// decision rather than an accidental leak in some future request type.
809    pub fn log_summary(&self) -> String {
810        match self {
811            AgentRequest::Ping => "Ping".into(),
812            AgentRequest::FsNotify { events } => format!("FsNotify {{ count: {} }}", events.len()),
813            AgentRequest::Pull { image, .. } => format!("Pull {{ image: {image} }}"),
814            AgentRequest::Query { image, .. } => format!("Query {{ image: {image} }}"),
815            AgentRequest::ListImages => "ListImages".into(),
816            AgentRequest::GarbageCollect { .. } => "GarbageCollect".into(),
817            AgentRequest::PrepareOverlay { .. } => "PrepareOverlay".into(),
818            AgentRequest::CleanupOverlay { .. } => "CleanupOverlay".into(),
819            AgentRequest::FormatStorage => "FormatStorage".into(),
820            AgentRequest::StorageStatus => "StorageStatus".into(),
821            AgentRequest::MemoryStatus => "MemoryStatus".into(),
822            AgentRequest::NetworkTest { .. } => "NetworkTest".into(),
823            AgentRequest::Shutdown { .. } => "Shutdown".into(),
824            AgentRequest::ExportLayer { .. } => "ExportLayer".into(),
825            AgentRequest::FlattenLayers { lowerdirs, .. } => {
826                format!("FlattenLayers {{ count: {} }}", lowerdirs.len())
827            }
828            AgentRequest::VmExec { .. } => "VmExec".into(),
829            AgentRequest::BranchpointWait { .. } => "BranchpointWait".into(),
830            AgentRequest::BranchpointArm => "BranchpointArm".into(),
831            AgentRequest::BranchpointPark => "BranchpointPark".into(),
832            AgentRequest::BranchpointRelease { .. } => "BranchpointRelease".into(),
833            AgentRequest::BranchpointActivate { .. } => "BranchpointActivate".into(),
834            AgentRequest::BranchpointWaitWorkerReady { .. } => "BranchpointWaitWorkerReady".into(),
835            AgentRequest::Run { image, .. } => format!("Run {{ image: {image} }}"),
836            AgentRequest::Stdin { .. } => "Stdin".into(),
837            AgentRequest::Resize { .. } => "Resize".into(),
838            AgentRequest::FileWrite { .. } => "FileWrite".into(),
839            AgentRequest::FileWriteBegin { .. } => "FileWriteBegin".into(),
840            AgentRequest::FileWriteChunk { .. } => "FileWriteChunk".into(),
841            AgentRequest::FileRead { .. } => "FileRead".into(),
842            AgentRequest::ListDirectory { .. } => "ListDirectory".into(),
843            AgentRequest::ArchiveDirectory { .. } => "ArchiveDirectory".into(),
844            // Pod requests: spec/process JSON may carry env secrets — emit ids only.
845            AgentRequest::PodCreate { id, .. } => format!("PodCreate {{ id: {id} }}"),
846            AgentRequest::PodStart { id, exec_id } => match exec_id {
847                Some(e) => format!("PodStart {{ id: {id}, exec: {e} }}"),
848                None => format!("PodStart {{ id: {id} }}"),
849            },
850            AgentRequest::PodExec { id, exec_id, .. } => {
851                format!("PodExec {{ id: {id}, exec: {exec_id} }}")
852            }
853            AgentRequest::PodSignal {
854                id, signal, all, ..
855            } => format!("PodSignal {{ id: {id}, signal: {signal}, all: {all} }}"),
856            AgentRequest::PodPids { id } => format!("PodPids {{ id: {id} }}"),
857            AgentRequest::PodStats { id } => format!("PodStats {{ id: {id} }}"),
858            AgentRequest::PodDelete { id, exec_id } => match exec_id {
859                Some(e) => format!("PodDelete {{ id: {id}, exec: {e} }}"),
860                None => format!("PodDelete {{ id: {id} }}"),
861            },
862        }
863    }
864}
865
866/// One entry of a [`AgentRequest::ListDirectory`] result.
867///
868/// Deliberately small: a name, what it is, and a size. A caller discovering a
869/// tree wants to know where to recurse and what is worth downloading, and
870/// anything richer invites the listing being used as a stat cache.
871#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
872pub struct DirectoryEntry {
873    /// Entry name, not a path: no separators, relative to the directory listed.
874    pub name: String,
875    /// `"file"`, `"dir"`, `"symlink"`, or `"other"` for sockets, devices and
876    /// fifos, which a caller can see but should not try to download.
877    pub kind: String,
878    /// Size in bytes. Zero for anything that is not a regular file.
879    #[serde(default)]
880    pub size: u64,
881}
882
883/// Agent response types.
884#[derive(Debug, Clone, Serialize, Deserialize)]
885#[serde(tag = "status", rename_all = "snake_case")]
886pub enum AgentResponse {
887    /// Operation completed successfully.
888    Ok {
889        /// Response data (varies by request type).
890        #[serde(default, skip_serializing_if = "Option::is_none")]
891        data: Option<serde_json::Value>,
892    },
893
894    /// Pong response to ping.
895    Pong {
896        /// Protocol version.
897        version: u32,
898        /// Optional agent features that can evolve independently of the base protocol.
899        #[serde(default, skip_serializing_if = "Vec::is_empty")]
900        capabilities: Vec<String>,
901    },
902
903    /// Progress update (for long operations like pull).
904    Progress {
905        /// Human-readable message.
906        message: String,
907        /// Completion percentage (0-100).
908        #[serde(default, skip_serializing_if = "Option::is_none")]
909        percent: Option<u8>,
910        /// Current layer being processed.
911        #[serde(default, skip_serializing_if = "Option::is_none")]
912        layer: Option<String>,
913    },
914
915    /// Operation failed.
916    Error {
917        /// Error message.
918        message: String,
919        /// Error code (for programmatic handling).
920        #[serde(default, skip_serializing_if = "Option::is_none")]
921        code: Option<String>,
922    },
923
924    /// Command execution completed (non-interactive mode).
925    Completed {
926        /// Exit code from the command.
927        exit_code: i32,
928        /// Standard output (may be truncated). `Vec<u8>` preserves binary
929        /// output (image bytes, tarballs, etc.) that would be truncated by
930        /// `String` at the first non-UTF-8 byte. Serialized as base64 JSON
931        /// string — the same format as the streaming `Stdout` variant.
932        #[serde(with = "base64_bytes")]
933        stdout: Vec<u8>,
934        /// Standard error (may be truncated).
935        #[serde(with = "base64_bytes")]
936        stderr: Vec<u8>,
937    },
938
939    /// Command started (interactive mode).
940    /// Indicates the command is running and ready to receive stdin.
941    Started,
942
943    /// Stdout data from a running command (interactive mode).
944    Stdout {
945        /// Output data.
946        #[serde(with = "base64_bytes")]
947        data: Vec<u8>,
948    },
949
950    /// Stderr data from a running command (interactive mode).
951    Stderr {
952        /// Error output data.
953        #[serde(with = "base64_bytes")]
954        data: Vec<u8>,
955    },
956
957    /// Command exited (interactive mode).
958    Exited {
959        /// Exit code from the command.
960        exit_code: i32,
961        /// The container was terminated by the cgroup OOM killer. The shim
962        /// turns this into a TaskOOM event so the CRI reports
963        /// `reason=OOMKilled`. Only ever set on a pod container's init exit.
964        #[serde(default)]
965        oom: bool,
966    },
967
968    /// PIDs inside a pod container (`PodPids` reply).
969    Pids {
970        /// Guest PIDs, container-init first when known.
971        pids: Vec<u32>,
972    },
973
974    /// Resource usage sample for a pod container (`PodStats` reply). Summed over
975    /// the container's process tree read from /proc (no per-container cgroup).
976    Stats {
977        /// Cumulative CPU time of the process tree, in nanoseconds.
978        cpu_usage_ns: u64,
979        /// Resident memory of the process tree, in bytes.
980        memory_bytes: u64,
981    },
982
983    /// Streaming binary-data chunk.
984    ///
985    /// Used by every streaming download direction: the agent sends
986    /// one or more `DataChunk` responses in sequence, with `done: true`
987    /// on the final chunk. Current producers: `ExportLayer` and
988    /// `FileRead`.
989    ///
990    /// Payload size per chunk should stay under
991    /// [`LAYER_CHUNK_SIZE`] so the encoded frame (~1.33× after
992    /// base64) fits inside [`MAX_FRAME_SIZE`] with JSON overhead to
993    /// spare.
994    DataChunk {
995        /// Chunk bytes. Empty allowed on the final frame (common for
996        /// EOF-on-clean-boundary cases).
997        #[serde(with = "base64_bytes")]
998        data: Vec<u8>,
999        /// True on the final chunk of the stream.
1000        done: bool,
1001    },
1002}
1003
1004// ============================================================================
1005// Error Code Constants
1006// ============================================================================
1007//
1008// Standard error codes for AgentResponse::Error. Using constants ensures
1009// consistency across the codebase and makes error handling more reliable.
1010
1011/// Error codes for agent responses.
1012pub mod error_codes {
1013    /// Request payload was invalid or malformed.
1014    pub const INVALID_REQUEST: &str = "INVALID_REQUEST";
1015    /// Requested resource was not found.
1016    pub const NOT_FOUND: &str = "NOT_FOUND";
1017    /// Internal error during operation.
1018    pub const INTERNAL_ERROR: &str = "INTERNAL_ERROR";
1019    /// Image pull operation failed.
1020    pub const PULL_FAILED: &str = "PULL_FAILED";
1021    /// Image query operation failed.
1022    pub const QUERY_FAILED: &str = "QUERY_FAILED";
1023    /// Command execution failed.
1024    pub const RUN_FAILED: &str = "RUN_FAILED";
1025    /// Command execution failed in container.
1026    pub const EXEC_FAILED: &str = "EXEC_FAILED";
1027    /// Process spawn failed.
1028    pub const SPAWN_FAILED: &str = "SPAWN_FAILED";
1029    /// Mount operation failed.
1030    pub const MOUNT_FAILED: &str = "MOUNT_FAILED";
1031    /// File I/O operation failed.
1032    pub const FILE_IO_FAILED: &str = "FILE_IO_FAILED";
1033    /// Overlay filesystem operation failed.
1034    pub const OVERLAY_FAILED: &str = "OVERLAY_FAILED";
1035    /// Cleanup operation failed.
1036    pub const CLEANUP_FAILED: &str = "CLEANUP_FAILED";
1037    /// Storage format operation failed.
1038    pub const FORMAT_FAILED: &str = "FORMAT_FAILED";
1039    /// Storage status query failed.
1040    pub const STATUS_FAILED: &str = "STATUS_FAILED";
1041    /// List operation failed.
1042    pub const LIST_FAILED: &str = "LIST_FAILED";
1043    /// Garbage collection failed.
1044    pub const GC_FAILED: &str = "GC_FAILED";
1045    /// Container creation failed.
1046    pub const CREATE_FAILED: &str = "CREATE_FAILED";
1047    /// Container start failed.
1048    pub const START_FAILED: &str = "START_FAILED";
1049    /// Container stop failed.
1050    pub const STOP_FAILED: &str = "STOP_FAILED";
1051    /// Container delete failed.
1052    pub const DELETE_FAILED: &str = "DELETE_FAILED";
1053    /// Export operation failed.
1054    pub const EXPORT_FAILED: &str = "EXPORT_FAILED";
1055    /// Serialization error.
1056    pub const SERIALIZATION_ERROR: &str = "SERIALIZATION_ERROR";
1057    /// Message size exceeds maximum.
1058    pub const MESSAGE_TOO_LARGE: &str = "MESSAGE_TOO_LARGE";
1059    /// Process wait operation failed.
1060    pub const WAIT_FAILED: &str = "WAIT_FAILED";
1061}
1062
1063impl AgentResponse {
1064    /// Create an error response with the given message and code.
1065    ///
1066    /// # Example
1067    ///
1068    /// ```
1069    /// use smolvm_protocol::{AgentResponse, error_codes};
1070    ///
1071    /// let response = AgentResponse::error("image not found", error_codes::NOT_FOUND);
1072    /// ```
1073    pub fn error(message: impl Into<String>, code: &str) -> Self {
1074        AgentResponse::Error {
1075            message: message.into(),
1076            code: Some(code.to_string()),
1077        }
1078    }
1079
1080    /// Create an error response from a Result's error, with the given code.
1081    ///
1082    /// # Example
1083    ///
1084    /// ```ignore
1085    /// let response = some_operation()
1086    ///     .map(|data| AgentResponse::ok_with_data(data))
1087    ///     .unwrap_or_else(|e| AgentResponse::from_err(e, error_codes::PULL_FAILED));
1088    /// ```
1089    pub fn from_err<E: std::fmt::Display>(err: E, code: &str) -> Self {
1090        AgentResponse::Error {
1091            message: err.to_string(),
1092            code: Some(code.to_string()),
1093        }
1094    }
1095
1096    /// Create an Ok response with optional JSON data.
1097    pub fn ok(data: Option<serde_json::Value>) -> Self {
1098        AgentResponse::Ok { data }
1099    }
1100
1101    /// Create an Ok response with JSON-serializable data.
1102    ///
1103    /// Returns an error response if serialization fails.
1104    pub fn ok_with_data<T: serde::Serialize>(data: T) -> Self {
1105        match serde_json::to_value(data) {
1106            Ok(value) => AgentResponse::Ok { data: Some(value) },
1107            Err(e) => AgentResponse::error(
1108                format!("failed to serialize response: {}", e),
1109                error_codes::SERIALIZATION_ERROR,
1110            ),
1111        }
1112    }
1113
1114    /// Convert a Result into an AgentResponse.
1115    ///
1116    /// On success, serializes the value to JSON. On error, creates an error response.
1117    ///
1118    /// # Example
1119    ///
1120    /// ```ignore
1121    /// let response = AgentResponse::from_result(
1122    ///     storage::pull_image(image),
1123    ///     error_codes::PULL_FAILED,
1124    /// );
1125    /// ```
1126    pub fn from_result<T, E>(result: Result<T, E>, error_code: &str) -> Self
1127    where
1128        T: serde::Serialize,
1129        E: std::fmt::Display,
1130    {
1131        match result {
1132            Ok(data) => Self::ok_with_data(data),
1133            Err(e) => Self::from_err(e, error_code),
1134        }
1135    }
1136}
1137
1138/// Image information returned by Query/ListImages.
1139#[derive(Debug, Clone, Serialize, Deserialize)]
1140pub struct ImageInfo {
1141    /// Image reference.
1142    pub reference: String,
1143    /// Image digest (sha256:...).
1144    pub digest: String,
1145    /// Image size in bytes.
1146    pub size: u64,
1147    /// Creation timestamp (ISO 8601).
1148    pub created: Option<String>,
1149    /// Platform architecture.
1150    pub architecture: String,
1151    /// Platform OS.
1152    pub os: String,
1153    /// Number of layers.
1154    pub layer_count: usize,
1155    /// Layer digests in order.
1156    pub layers: Vec<String>,
1157    /// Image entrypoint (from OCI config).
1158    #[serde(default)]
1159    pub entrypoint: Vec<String>,
1160    /// Image default command (from OCI config).
1161    #[serde(default)]
1162    pub cmd: Vec<String>,
1163    /// Image environment variables (from OCI config).
1164    #[serde(default)]
1165    pub env: Vec<String>,
1166    /// Image working directory (from OCI config).
1167    #[serde(default)]
1168    pub workdir: Option<String>,
1169    /// Image default user (from OCI config).
1170    #[serde(default)]
1171    pub user: Option<String>,
1172}
1173
1174/// Overlay preparation result.
1175#[derive(Debug, Clone, Serialize, Deserialize)]
1176pub struct OverlayInfo {
1177    /// Path to the merged overlay rootfs.
1178    pub rootfs_path: String,
1179    /// Path to the upper (writable) directory.
1180    pub upper_path: String,
1181    /// Path to the work directory.
1182    pub work_path: String,
1183}
1184
1185/// The guest's own account of machine memory, from `/proc/meminfo`.
1186///
1187/// This is what a machine is actually using. The host-side `phys_footprint`
1188/// answers a different question — it counts compressed pages at their
1189/// uncompressed size and keeps counting memory the guest has stopped needing —
1190/// so it runs well above these figures and should not be read as consumption.
1191#[derive(Debug, Clone, Copy, Default, Serialize, Deserialize, PartialEq, Eq)]
1192#[serde(rename_all = "camelCase")]
1193pub struct MemoryStatus {
1194    /// Total usable RAM the guest sees, in bytes. Slightly below the machine's
1195    /// configured size: the kernel reserves some before the allocator sees it.
1196    pub total_bytes: u64,
1197    /// Memory available for new allocations without swapping, in bytes. The
1198    /// figure to report, since `free_bytes` excludes reclaimable page cache.
1199    pub available_bytes: u64,
1200    /// Free memory in bytes: never allocated, or released and not reused.
1201    pub free_bytes: u64,
1202    /// Page cache in bytes. Counted inside `available_bytes`.
1203    pub cached_bytes: u64,
1204    /// Swap configured in the guest, in bytes. Zero when the guest has none.
1205    pub swap_total_bytes: u64,
1206    /// Swap in use, in bytes.
1207    pub swap_used_bytes: u64,
1208}
1209
1210impl MemoryStatus {
1211    /// Memory the guest cannot hand back on demand.
1212    pub fn used_bytes(&self) -> u64 {
1213        self.total_bytes.saturating_sub(self.available_bytes)
1214    }
1215}
1216
1217/// Storage status information.
1218#[derive(Debug, Clone, Serialize, Deserialize)]
1219pub struct StorageStatus {
1220    /// Whether the storage is formatted and ready.
1221    pub ready: bool,
1222    /// Total size in bytes.
1223    pub total_bytes: u64,
1224    /// Used size in bytes.
1225    pub used_bytes: u64,
1226    /// Number of cached layers.
1227    pub layer_count: usize,
1228    /// Number of cached images.
1229    pub image_count: usize,
1230}
1231
1232/// Registry authentication credentials for pulling images.
1233///
1234/// `Debug` is hand-written to redact the password: this value is carried inside
1235/// `AgentRequest::Pull`, and any `{:?}` of that request (e.g. a tracing span)
1236/// would otherwise serialize the token verbatim into the machine's console log,
1237/// which is exposed over the logs API.
1238#[derive(Clone, Serialize, Deserialize)]
1239pub struct RegistryAuth {
1240    /// Username for authentication.
1241    pub username: String,
1242    /// Password or token for authentication.
1243    pub password: String,
1244}
1245
1246impl std::fmt::Debug for RegistryAuth {
1247    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1248        f.debug_struct("RegistryAuth")
1249            .field("username", &self.username)
1250            .field("password", &"***")
1251            .finish()
1252    }
1253}
1254
1255// ============================================================================
1256// Workload VM Protocol (Command Execution)
1257// ============================================================================
1258
1259/// Messages from host to workload VM.
1260#[derive(Debug, Clone, Serialize, Deserialize)]
1261#[serde(tag = "type", rename_all = "snake_case")]
1262pub enum HostMessage {
1263    /// Authentication request.
1264    Auth {
1265        /// Authentication token (base64).
1266        token: String,
1267        /// Protocol version.
1268        protocol_version: u32,
1269    },
1270
1271    /// Run a command.
1272    Run {
1273        /// Request ID for correlating responses.
1274        request_id: u64,
1275        /// Command and arguments.
1276        command: Vec<String>,
1277        /// Environment variables.
1278        env: Vec<(String, String)>,
1279        /// Working directory.
1280        workdir: Option<String>,
1281    },
1282
1283    /// Execute a command in running VM.
1284    Exec {
1285        /// Request ID.
1286        request_id: u64,
1287        /// Command and arguments.
1288        command: Vec<String>,
1289        /// Allocate a TTY.
1290        tty: bool,
1291    },
1292
1293    /// Send a signal to a running command.
1294    Signal {
1295        /// Request ID of the command.
1296        request_id: u64,
1297        /// Signal number.
1298        signal: i32,
1299    },
1300
1301    /// Request graceful shutdown.
1302    Stop {
1303        /// Timeout in milliseconds.
1304        timeout_ms: u64,
1305    },
1306}
1307
1308/// Messages from workload VM to host.
1309#[derive(Debug, Clone, Serialize, Deserialize)]
1310#[serde(tag = "type", rename_all = "snake_case")]
1311pub enum GuestMessage {
1312    /// Authentication successful.
1313    AuthOk,
1314
1315    /// Authentication failed.
1316    AuthFailed,
1317
1318    /// VM is ready to receive commands.
1319    Ready,
1320
1321    /// Command started.
1322    Started {
1323        /// Request ID.
1324        request_id: u64,
1325    },
1326
1327    /// Stdout data from command.
1328    Stdout {
1329        /// Request ID.
1330        request_id: u64,
1331        /// Output data.
1332        #[serde(with = "base64_bytes")]
1333        data: Vec<u8>,
1334        /// Whether output was truncated.
1335        truncated: bool,
1336    },
1337
1338    /// Stderr data from command.
1339    Stderr {
1340        /// Request ID.
1341        request_id: u64,
1342        /// Output data.
1343        #[serde(with = "base64_bytes")]
1344        data: Vec<u8>,
1345        /// Whether output was truncated.
1346        truncated: bool,
1347    },
1348
1349    /// Command exited.
1350    Exit {
1351        /// Request ID.
1352        request_id: u64,
1353        /// Exit code.
1354        code: i32,
1355        /// Exit reason.
1356        reason: String,
1357    },
1358
1359    /// Error occurred.
1360    Error {
1361        /// Request ID (if applicable).
1362        request_id: Option<u64>,
1363        /// Error message.
1364        message: String,
1365    },
1366}
1367
1368// ============================================================================
1369// Wire Format Helpers
1370// ============================================================================
1371
1372/// Envelope that wraps any message with an optional trace ID for correlation.
1373///
1374/// On the wire, the trace_id is flattened into the JSON alongside the message
1375/// fields: `{"trace_id":"abc123","method":"ping"}`.
1376#[derive(Debug, Clone, Serialize, Deserialize)]
1377pub struct Envelope<T> {
1378    /// Trace ID for correlating host API requests to agent operations.
1379    #[serde(skip_serializing_if = "Option::is_none", default)]
1380    pub trace_id: Option<String>,
1381    /// The wrapped message.
1382    #[serde(flatten)]
1383    pub body: T,
1384}
1385
1386impl<T> Envelope<T> {
1387    /// Create an envelope with no trace ID.
1388    pub fn new(body: T) -> Self {
1389        Self {
1390            trace_id: None,
1391            body,
1392        }
1393    }
1394
1395    /// Create an envelope with an optional trace ID.
1396    pub fn with_trace_id(body: T, trace_id: Option<String>) -> Self {
1397        Self { trace_id, body }
1398    }
1399}
1400
1401/// Encode a message to wire format (length-prefixed JSON).
1402pub fn encode_message<T: Serialize>(msg: &T) -> Result<Vec<u8>, serde_json::Error> {
1403    let json = serde_json::to_vec(msg)?;
1404    let len = json.len() as u32;
1405
1406    let mut buf = Vec::with_capacity(4 + json.len());
1407    buf.extend_from_slice(&len.to_be_bytes());
1408    buf.extend_from_slice(&json);
1409
1410    Ok(buf)
1411}
1412
1413/// Decode a message from wire format.
1414pub fn decode_message<T: for<'de> Deserialize<'de>>(data: &[u8]) -> Result<T, DecodeError> {
1415    if data.len() < 4 {
1416        return Err(DecodeError::TooShort);
1417    }
1418
1419    let len = u32::from_be_bytes([data[0], data[1], data[2], data[3]]) as usize;
1420
1421    if len > MAX_FRAME_SIZE as usize {
1422        return Err(DecodeError::TooLarge(len));
1423    }
1424
1425    if data.len() < 4 + len {
1426        return Err(DecodeError::Incomplete {
1427            expected: len,
1428            got: data.len() - 4,
1429        });
1430    }
1431
1432    serde_json::from_slice(&data[4..4 + len]).map_err(DecodeError::Json)
1433}
1434
1435/// Error decoding a wire message.
1436#[derive(Debug)]
1437pub enum DecodeError {
1438    /// Data too short to contain length header.
1439    TooShort,
1440    /// Frame size exceeds maximum.
1441    TooLarge(usize),
1442    /// Incomplete frame.
1443    Incomplete {
1444        /// Expected length.
1445        expected: usize,
1446        /// Actual length.
1447        got: usize,
1448    },
1449    /// JSON parse error.
1450    Json(serde_json::Error),
1451}
1452
1453impl std::fmt::Display for DecodeError {
1454    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1455        match self {
1456            DecodeError::TooShort => write!(f, "data too short for length header"),
1457            DecodeError::TooLarge(size) => write!(f, "frame too large: {} bytes", size),
1458            DecodeError::Incomplete { expected, got } => {
1459                write!(
1460                    f,
1461                    "incomplete frame: expected {} bytes, got {}",
1462                    expected, got
1463                )
1464            }
1465            DecodeError::Json(e) => write!(f, "JSON decode error: {}", e),
1466        }
1467    }
1468}
1469
1470impl std::error::Error for DecodeError {}
1471
1472#[cfg(test)]
1473mod tests {
1474    use super::*;
1475
1476    #[test]
1477    fn file_write_without_owner_fields_still_parses() {
1478        // Requests from clients predating uid/gid must keep deserializing.
1479        let old = r#"{"method":"file_write","path":"/x","data":"aGk=","mode":420}"#;
1480        let req: AgentRequest = serde_json::from_str(old).unwrap();
1481        match req {
1482            AgentRequest::FileWrite { mode, uid, gid, .. } => {
1483                assert_eq!(mode, Some(420));
1484                assert_eq!(uid, None);
1485                assert_eq!(gid, None);
1486            }
1487            other => panic!("unexpected: {other:?}"),
1488        }
1489    }
1490
1491    #[test]
1492    fn flatten_layers_output_is_optional() {
1493        // A client predating streaming always names a guest path to stage the
1494        // archive at, and must keep deserializing.
1495        let staged =
1496            r#"{"method":"flatten_layers","lowerdirs":["/a","/b"],"output":"/storage/x.tar"}"#;
1497        let req: AgentRequest = serde_json::from_str(staged).unwrap();
1498        match req {
1499            AgentRequest::FlattenLayers { lowerdirs, output } => {
1500                assert_eq!(lowerdirs, vec!["/a".to_string(), "/b".to_string()]);
1501                assert_eq!(output, Some("/storage/x.tar".to_string()));
1502            }
1503            other => panic!("unexpected: {other:?}"),
1504        }
1505
1506        // Omitting it asks the agent to stream the archive back instead, which is
1507        // what keeps a large export from needing room for a second full copy.
1508        let streamed = r#"{"method":"flatten_layers","lowerdirs":["/a","/b"]}"#;
1509        let req: AgentRequest = serde_json::from_str(streamed).unwrap();
1510        match req {
1511            AgentRequest::FlattenLayers { output, .. } => assert_eq!(output, None),
1512            other => panic!("unexpected: {other:?}"),
1513        }
1514    }
1515
1516    #[test]
1517    fn test_encode_decode_roundtrip() {
1518        let req = AgentRequest::Pull {
1519            image: "alpine:latest".to_string(),
1520            oci_platform: Some("linux/arm64".to_string()),
1521            auth: None,
1522            proxy: None,
1523            no_proxy: None,
1524        };
1525
1526        let encoded = encode_message(&req).unwrap();
1527        let decoded: AgentRequest = decode_message(&encoded).unwrap();
1528
1529        let AgentRequest::Pull {
1530            image,
1531            oci_platform,
1532            auth,
1533            proxy,
1534            no_proxy,
1535        } = decoded
1536        else {
1537            panic!("expected Pull variant, got {:?}", decoded);
1538        };
1539        assert_eq!(image, "alpine:latest");
1540        assert_eq!(oci_platform, Some("linux/arm64".to_string()));
1541        assert!(auth.is_none());
1542        assert!(proxy.is_none());
1543        assert!(no_proxy.is_none());
1544    }
1545
1546    #[test]
1547    fn test_encode_decode_with_auth() {
1548        let req = AgentRequest::Pull {
1549            image: "ghcr.io/owner/repo:latest".to_string(),
1550            oci_platform: None,
1551            auth: Some(RegistryAuth {
1552                username: "testuser".to_string(),
1553                password: "testpass".to_string(),
1554            }),
1555            proxy: None,
1556            no_proxy: None,
1557        };
1558
1559        let encoded = encode_message(&req).unwrap();
1560        let decoded: AgentRequest = decode_message(&encoded).unwrap();
1561
1562        let AgentRequest::Pull {
1563            image,
1564            oci_platform,
1565            auth,
1566            proxy: _,
1567            no_proxy: _,
1568        } = decoded
1569        else {
1570            panic!("expected Pull variant, got {:?}", decoded);
1571        };
1572        assert_eq!(image, "ghcr.io/owner/repo:latest");
1573        assert!(oci_platform.is_none());
1574        let auth = auth.expect("auth should be Some");
1575        assert_eq!(auth.username, "testuser");
1576        assert_eq!(auth.password, "testpass");
1577    }
1578
1579    #[test]
1580    fn test_encode_decode_with_proxy() {
1581        let req = AgentRequest::Pull {
1582            image: "alpine:latest".to_string(),
1583            oci_platform: None,
1584            auth: None,
1585            proxy: Some("http://192.168.127.254:3128".to_string()),
1586            no_proxy: Some("127.0.0.1,localhost,.internal".to_string()),
1587        };
1588
1589        let encoded = encode_message(&req).unwrap();
1590        let decoded: AgentRequest = decode_message(&encoded).unwrap();
1591
1592        let AgentRequest::Pull {
1593            proxy, no_proxy, ..
1594        } = decoded
1595        else {
1596            panic!("expected Pull variant, got {:?}", decoded);
1597        };
1598        assert_eq!(proxy.as_deref(), Some("http://192.168.127.254:3128"));
1599        assert_eq!(no_proxy.as_deref(), Some("127.0.0.1,localhost,.internal"));
1600    }
1601
1602    #[test]
1603    fn test_decode_too_short() {
1604        let data = [0u8; 2];
1605        let result: Result<AgentRequest, _> = decode_message(&data);
1606        assert!(matches!(result, Err(DecodeError::TooShort)));
1607    }
1608
1609    #[test]
1610    fn test_decode_incomplete() {
1611        let mut data = vec![0, 0, 0, 100]; // claims 100 bytes
1612        data.extend_from_slice(b"{}"); // only 2 bytes of payload
1613        let result: Result<AgentRequest, _> = decode_message(&data);
1614        assert!(matches!(result, Err(DecodeError::Incomplete { .. })));
1615    }
1616
1617    #[test]
1618    fn test_agent_request_serialization() {
1619        let req = AgentRequest::Ping;
1620        let json = serde_json::to_string(&req).unwrap();
1621        assert!(json.contains("ping"));
1622
1623        let req = AgentRequest::PrepareOverlay {
1624            image: "ubuntu:22.04".to_string(),
1625            workload_id: "wl-123".to_string(),
1626        };
1627        let json = serde_json::to_string(&req).unwrap();
1628        assert!(json.contains("prepare_overlay"));
1629    }
1630
1631    #[test]
1632    fn test_agent_response_serialization() {
1633        let resp = AgentResponse::Pong {
1634            version: PROTOCOL_VERSION,
1635            capabilities: vec![forkpoint::TYPED_BRANCHPOINT_CAPABILITY.to_string()],
1636        };
1637        let json = serde_json::to_string(&resp).unwrap();
1638        assert!(json.contains("pong"));
1639        assert!(json.contains(forkpoint::TYPED_BRANCHPOINT_CAPABILITY));
1640
1641        let legacy: AgentResponse =
1642            serde_json::from_str(r#"{"status":"pong","version":1}"#).unwrap();
1643        assert!(matches!(
1644            legacy,
1645            AgentResponse::Pong {
1646                version: PROTOCOL_VERSION,
1647                capabilities
1648            } if capabilities.is_empty()
1649        ));
1650
1651        let resp = AgentResponse::Progress {
1652            message: "Pulling layer 1/3".to_string(),
1653            percent: Some(33),
1654            layer: Some("sha256:abc123".to_string()),
1655        };
1656        let json = serde_json::to_string(&resp).unwrap();
1657        assert!(json.contains("progress"));
1658    }
1659
1660    #[test]
1661    fn file_write_begin_roundtrips() {
1662        let req = AgentRequest::FileWriteBegin {
1663            path: "/tmp/target".into(),
1664            mode: Some(0o600),
1665            uid: Some(1000),
1666            gid: Some(1000),
1667            total_size: 123_456_789,
1668        };
1669        let bytes = encode_message(&req).unwrap();
1670        let back: AgentRequest = decode_message(&bytes).unwrap();
1671        match back {
1672            AgentRequest::FileWriteBegin {
1673                path,
1674                mode,
1675                uid,
1676                gid,
1677                total_size,
1678            } => {
1679                assert_eq!(path, "/tmp/target");
1680                assert_eq!(mode, Some(0o600));
1681                assert_eq!(uid, Some(1000));
1682                assert_eq!(gid, Some(1000));
1683                assert_eq!(total_size, 123_456_789);
1684            }
1685            _ => panic!("wrong variant"),
1686        }
1687    }
1688
1689    #[test]
1690    fn file_write_chunk_roundtrips_binary_data() {
1691        // Binary data (bytes outside UTF-8) must survive the base64
1692        // trip intact. If the encoding ever silently lossifies, this
1693        // fires.
1694        let payload: Vec<u8> = (0u8..=255).collect();
1695        let req = AgentRequest::FileWriteChunk {
1696            data: payload.clone(),
1697            done: true,
1698        };
1699        let bytes = encode_message(&req).unwrap();
1700        let back: AgentRequest = decode_message(&bytes).unwrap();
1701        match back {
1702            AgentRequest::FileWriteChunk { data, done } => {
1703                assert_eq!(data, payload);
1704                assert!(done);
1705            }
1706            _ => panic!("wrong variant"),
1707        }
1708    }
1709
1710    #[test]
1711    fn file_write_size_constants_are_frame_safe() {
1712        // Sanity: a single streaming chunk at FILE_WRITE_CHUNK_SIZE
1713        // must fit inside MAX_FRAME_SIZE after base64 (+ ~33%) and
1714        // JSON overhead. If anyone bumps CHUNK_SIZE past the limit,
1715        // this test fires before production does.
1716        let chunk_bytes = FILE_WRITE_CHUNK_SIZE as u64;
1717        let base64_bytes = chunk_bytes.div_ceil(3) * 4; // ceil(n/3)*4
1718        let json_overhead = 256u64; // method tag, done bool, quotes
1719        let total = base64_bytes + json_overhead;
1720        assert!(
1721            total < MAX_FRAME_SIZE as u64,
1722            "FILE_WRITE_CHUNK_SIZE of {} bytes would produce a frame \
1723             of ~{} bytes which exceeds MAX_FRAME_SIZE of {}",
1724            chunk_bytes,
1725            total,
1726            MAX_FRAME_SIZE
1727        );
1728    }
1729
1730    #[test]
1731    fn test_ports_constants() {
1732        assert_eq!(ports::WORKLOAD_CONTROL, 5000);
1733        assert_eq!(ports::WORKLOAD_LOGS, 5001);
1734        assert_eq!(ports::AGENT_CONTROL, 6000);
1735        assert_eq!(ports::SSH_AGENT, 6001);
1736    }
1737
1738    #[test]
1739    fn test_cid_constants() {
1740        assert_eq!(cid::HOST, 2);
1741        assert_eq!(cid::GUEST, 3);
1742    }
1743
1744    #[test]
1745    fn test_envelope_serialization_with_trace_id() {
1746        let req = AgentRequest::Ping;
1747        let envelope = Envelope::with_trace_id(&req, Some("abc123".to_string()));
1748        let json = serde_json::to_string(&envelope).unwrap();
1749
1750        // trace_id should be flattened alongside the method tag
1751        assert!(json.contains("\"trace_id\":\"abc123\""));
1752        assert!(json.contains("\"method\":\"ping\""));
1753
1754        // Deserialize back — Envelope<AgentRequest> with flatten
1755        let parsed: Envelope<AgentRequest> = serde_json::from_str(&json).unwrap();
1756        assert_eq!(parsed.trace_id.as_deref(), Some("abc123"));
1757        assert!(matches!(parsed.body, AgentRequest::Ping));
1758    }
1759
1760    #[test]
1761    fn test_envelope_without_trace_id() {
1762        let req = AgentRequest::Ping;
1763        let envelope = Envelope::new(&req);
1764        let json = serde_json::to_string(&envelope).unwrap();
1765
1766        // No trace_id field (skip_serializing_if = None)
1767        assert!(!json.contains("trace_id"));
1768        assert!(json.contains("\"method\":\"ping\""));
1769    }
1770
1771    #[test]
1772    fn test_envelope_backward_compat_bare_request() {
1773        // A bare AgentRequest (no Envelope) should fail to parse as Envelope
1774        // but succeed as bare AgentRequest — this is the agent's fallback path
1775        let bare_json = r#"{"method":"ping"}"#;
1776
1777        // Envelope parse should fail (no body field to flatten into)
1778        // Actually with flatten, this may work — let's verify
1779        let envelope_result = serde_json::from_str::<Envelope<AgentRequest>>(bare_json);
1780        let bare_result = serde_json::from_str::<AgentRequest>(bare_json);
1781
1782        // At least one must succeed for backward compat
1783        assert!(
1784            envelope_result.is_ok() || bare_result.is_ok(),
1785            "Neither Envelope nor bare parse succeeded"
1786        );
1787
1788        // Bare parse must always work
1789        assert!(bare_result.is_ok());
1790        assert!(matches!(bare_result.unwrap(), AgentRequest::Ping));
1791
1792        // If Envelope works, trace_id should be None
1793        if let Ok(env) = envelope_result {
1794            assert!(env.trace_id.is_none());
1795        }
1796    }
1797}