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