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