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    /// Test network connectivity directly from the agent (not via chroot).
337    /// Used to debug TSI networking.
338    NetworkTest {
339        /// URL to test (e.g., "http://1.1.1.1")
340        url: String,
341    },
342
343    /// Shutdown the agent.
344    Shutdown,
345
346    /// Export a layer as a tar archive.
347    ///
348    /// Used by `smolvm pack` to extract OCI layers for packaging.
349    /// The agent streams the layer tar data back via LayerData responses.
350    ExportLayer {
351        /// Image digest (sha256:...).
352        image_digest: String,
353        /// Layer index (0-based).
354        layer_index: usize,
355    },
356
357    /// Merge a stack of directories into a single tar archive.
358    ///
359    /// `pack create --from-vm` ships the machine's image layers plus whatever the
360    /// container has written since as ONE flattened layer. The merge has to apply
361    /// whiteouts and opaque markers exactly as the runtime would, so it is done
362    /// with a read-only overlay mount rather than a file copy.
363    ///
364    /// The agent owns this rather than the host driving `mount(8)` over `VmExec`,
365    /// because `mount(8)` rejects a `lowerdir=` value beyond ~255 bytes — about
366    /// three OCI layer paths — while the agent can append each layer separately
367    /// through the same `fsconfig` path the runtime container mount uses.
368    FlattenLayers {
369        /// Directories to merge, bottom -> top. Entries that are missing or empty
370        /// are dropped, so a caller may include a container overlay's upper dir
371        /// without first checking whether the machine ever wrote to it.
372        lowerdirs: Vec<String>,
373        /// Guest path to write the tar archive to.
374        output: String,
375    },
376
377    /// Wait until the workload has declared a branchpoint, returning the ready
378    /// marker's contents (its profile lines) in `data.contents`.
379    BranchpointWait {
380        /// Give up after this many milliseconds.
381        timeout_ms: u64,
382    },
383    /// Put a negotiated helper into its restore-safe loop before capture.
384    BranchpointArm,
385    /// Return a parked source to its ordinary wait after capture.
386    BranchpointPark,
387    /// Release a restored clone: write the release marker for the generation
388    /// recorded in its ready marker, carrying the clone's identity, and wait
389    /// for the helper to acknowledge.
390    BranchpointRelease {
391        /// The clone's parameters in dotenv form, appended to the marker.
392        #[serde(default)]
393        env_dotenv: Option<String>,
394    },
395    /// Assign and release a held clone in one idempotent step: claim it with
396    /// `activation_token` (a retry with the same token completes a partial
397    /// commit; a different token is refused), install the per-clone
398    /// parameters, then write the release marker.
399    BranchpointActivate {
400        /// Per-clone parameters, dotenv form, written to `env_path`.
401        env_dotenv: String,
402        /// The same parameters in shell-sourceable form, written to `branch_env_path`.
403        env_sourceable: String,
404        /// Guest path of the dotenv file (under the clone's overlay for image machines).
405        env_path: String,
406        /// Guest path of the sourceable file.
407        branch_env_path: String,
408        /// Directory that must already exist, typically the clone's merged
409        /// overlay root; activation refuses rather than fabricating it.
410        require_dir: Option<String>,
411        /// Directory to create before writing the env files.
412        env_dir: String,
413        /// Token identifying this activation attempt.
414        activation_token: String,
415    },
416    /// Wait until a released clone's workload publishes its worker-ready token.
417    BranchpointWaitWorkerReady {
418        /// The token the workload must publish; any other is a mismatch.
419        token: String,
420        /// Give up after this many milliseconds.
421        timeout_ms: u64,
422    },
423
424    /// Execute a command directly in the VM (not in a container).
425    ///
426    /// This runs the command in the agent's Alpine rootfs without any
427    /// container isolation. Useful for VM-level operations and debugging.
428    VmExec {
429        /// Command and arguments.
430        command: Vec<String>,
431        /// Environment variables.
432        #[serde(default)]
433        env: Vec<(String, String)>,
434        /// Working directory in the VM.
435        workdir: Option<String>,
436        /// Timeout in milliseconds.
437        #[serde(default)]
438        timeout_ms: Option<u64>,
439        /// Interactive mode - stream I/O instead of buffering.
440        #[serde(default)]
441        interactive: bool,
442        /// Allocate a pseudo-TTY for the command.
443        #[serde(default)]
444        tty: bool,
445        /// Background mode - spawn and return PID immediately without waiting.
446        #[serde(default)]
447        background: bool,
448        /// Data to pipe to the command's stdin.
449        #[serde(default)]
450        stdin_data: Option<String>,
451    },
452
453    /// Run a command in an image's rootfs.
454    ///
455    /// This prepares an overlay, chroots into it, and executes the command.
456    /// Returns stdout, stderr, and exit code when the command completes.
457    Run {
458        /// Image reference (must be pulled first).
459        image: String,
460        /// Command and arguments.
461        command: Vec<String>,
462        /// Environment variables.
463        #[serde(default)]
464        env: Vec<(String, String)>,
465        /// Working directory inside the rootfs.
466        workdir: Option<String>,
467        /// User inside the rootfs. If omitted, the OCI image default applies.
468        #[serde(default, skip_serializing_if = "Option::is_none")]
469        user: Option<String>,
470        /// Volume mounts to bind into the container.
471        /// Each tuple is (virtiofs_tag, container_path, read_only).
472        #[serde(default)]
473        mounts: Vec<(String, String, bool)>,
474        /// Timeout in milliseconds. If the command exceeds this duration,
475        /// it will be killed and return exit code 124.
476        #[serde(default)]
477        timeout_ms: Option<u64>,
478        /// Interactive mode - stream I/O instead of buffering.
479        /// When true, output is streamed via Stdout/Stderr responses,
480        /// and stdin can be sent via the Stdin request.
481        #[serde(default)]
482        interactive: bool,
483        /// Allocate a pseudo-TTY for the command.
484        /// Enables terminal features like colors, line editing, and signal handling.
485        #[serde(default)]
486        tty: bool,
487        /// Detached mode — start the container and return immediately with the
488        /// container ID. Only meaningful when `persistent_overlay_id` is set.
489        /// Returns a `Completed` response with `stdout` containing the container ID.
490        #[serde(default)]
491        detached: bool,
492        /// Run the workload as an unprivileged container: restricted capabilities,
493        /// read-only cgroup, and no extra tmpfs. The default (false) is "VM-grade"
494        /// — since the microVM is the isolation boundary, the workload gets a full
495        /// capability set and the mounts an init system needs (so any image, incl.
496        /// systemd, boots). Opt in for defense-in-depth when running untrusted code.
497        #[serde(default)]
498        unprivileged: bool,
499        /// If set, use a persistent overlay that survives across exec sessions.
500        /// The overlay is identified by this ID (typically the machine name)
501        /// and reused on subsequent runs. If not set, an ephemeral overlay is
502        /// created and destroyed after the run.
503        #[serde(default, skip_serializing_if = "Option::is_none")]
504        persistent_overlay_id: Option<String>,
505        /// Data to pipe to the command's stdin (non-interactive runs only).
506        /// The pipe is closed after writing, so the command sees EOF.
507        #[serde(default, skip_serializing_if = "Option::is_none")]
508        stdin_data: Option<String>,
509        /// Spawn the container and return immediately with the crun PID.
510        /// The container runs detached; stdout/stderr go to /dev/null.
511        /// Incompatible with `interactive` and `tty`.
512        #[serde(default)]
513        background: bool,
514        /// S3 volumes to mount into the workload container. The agent mounts
515        /// them between `crun create` and `crun start`, so the workload's first
516        /// instruction already sees them and its command is never rewritten.
517        #[serde(default, skip_serializing_if = "Vec::is_empty")]
518        s3_volumes: Vec<S3Volume>,
519    },
520
521    /// Send stdin data to a running interactive command.
522    Stdin {
523        /// Input data to send to the command's stdin.
524        #[serde(with = "base64_bytes")]
525        data: Vec<u8>,
526    },
527
528    /// Resize the PTY window (for TTY mode).
529    Resize {
530        /// New width in columns.
531        cols: u16,
532        /// New height in rows.
533        rows: u16,
534    },
535
536    // ========================================================================
537    // File I/O
538    // ========================================================================
539    /// Write a file inside the VM in a single message.
540    ///
541    /// Use only for files up to [`FILE_WRITE_SINGLE_SHOT_MAX`]. Larger
542    /// files must stream via [`Self::FileWriteBegin`] +
543    /// [`Self::FileWriteChunk`] to avoid exceeding [`MAX_FRAME_SIZE`]
544    /// after base64 + JSON inflation.
545    FileWrite {
546        /// Absolute path in the VM filesystem.
547        path: String,
548        /// File contents.
549        #[serde(with = "base64_bytes")]
550        data: Vec<u8>,
551        /// File mode (e.g., 0o644). None = default (0644).
552        #[serde(default)]
553        mode: Option<u32>,
554        /// Owner uid to apply after the write. None = leave as written (root).
555        #[serde(default)]
556        uid: Option<u32>,
557        /// Owner gid to apply after the write. None = leave as written (root).
558        #[serde(default)]
559        gid: Option<u32>,
560    },
561
562    /// Open a streaming file upload session on this connection.
563    ///
564    /// Must be followed by one or more [`Self::FileWriteChunk`]
565    /// requests. The final chunk sets `done: true` to finalize.
566    /// Dropping the connection (or sending any non-chunk request)
567    /// before `done` aborts the session and leaves no partial file
568    /// at `path`.
569    ///
570    /// Sessions are per-connection — one session at a time.
571    FileWriteBegin {
572        /// Absolute path in the VM filesystem.
573        path: String,
574        /// File mode (e.g., 0o644). None = default (0644).
575        #[serde(default)]
576        mode: Option<u32>,
577        /// Owner uid to apply on finalize. None = leave as written (root).
578        #[serde(default)]
579        uid: Option<u32>,
580        /// Owner gid to apply on finalize. None = leave as written (root).
581        #[serde(default)]
582        gid: Option<u32>,
583        /// Expected total size in bytes. Rejected if it exceeds
584        /// [`FILE_TRANSFER_MAX_TOTAL`]. The agent uses this for an
585        /// early-fail check only; the actual size written is the sum
586        /// of chunk byte lengths.
587        total_size: u64,
588    },
589
590    /// Append a chunk to the currently open streaming upload.
591    /// If `done` is true, the agent fsyncs and atomically renames the
592    /// staging file onto the target path.
593    FileWriteChunk {
594        /// Chunk bytes. Typically [`FILE_WRITE_CHUNK_SIZE`] except
595        /// for the last chunk.
596        #[serde(with = "base64_bytes")]
597        data: Vec<u8>,
598        /// True on the final chunk; closes and renames the staging
599        /// file. False on intermediate chunks.
600        done: bool,
601    },
602
603    /// Read a file from the VM.
604    FileRead {
605        /// Absolute path in the VM filesystem.
606        path: String,
607    },
608
609    /// Stream a tar archive of a guest directory without creating a guest-side
610    /// temporary file. Used by staged mounts to batch many small files across
611    /// vsock instead of paying one virtiofs round trip per file.
612    ArchiveDirectory {
613        /// Absolute directory path in the VM filesystem.
614        path: String,
615    },
616
617    /// Create (without starting) a Kubernetes pod container whose rootfs is a
618    /// virtiofs-shared host directory (containerd snapshotter output) and whose
619    /// process definition comes from the host's OCI config. The agent builds
620    /// its crun bundle around the shared rootfs; nothing runs until
621    /// `PodStart`. Part of the containerd shim v2 datapath
622    /// (docs/kubernetes-runtime.md).
623    PodCreate {
624        /// Container ID (containerd task id).
625        id: String,
626        /// Rootfs path relative to the sandbox's shared virtiofs mount. The
627        /// shim boots the sandbox VM with ONE shared dir and bind-mounts each
628        /// container's rootfs under it (virtiofs shares are fixed at boot, but
629        /// pod containers are created afterwards), so the guest resolves this
630        /// as `<sandbox-share-mount>/<rootfs_rel>`.
631        rootfs_rel: String,
632        /// The host OCI runtime spec (config.json bytes). The agent extracts
633        /// process/env/cwd/user/mounts/resources and grafts them onto its own
634        /// guest bundle template; host-specific namespaces/paths are ignored.
635        spec_json: String,
636        /// Allocate a PTY for the init process.
637        #[serde(default)]
638        tty: bool,
639    },
640
641    /// Start a pod container created by `PodCreate` (or an exec process
642    /// registered by `PodExec`), streaming its I/O on THIS connection:
643    /// `Started` → `Stdout`/`Stderr`... → `Exited`. Stdin arrives via `Stdin`
644    /// requests; PTY resize via `Resize`.
645    PodStart {
646        /// Container ID.
647        id: String,
648        /// Exec process to start instead of the init process.
649        #[serde(default, skip_serializing_if = "Option::is_none")]
650        exec_id: Option<String>,
651    },
652
653    /// Register an exec process for a running pod container. Started later by
654    /// `PodStart { exec_id }`.
655    PodExec {
656        /// Container ID.
657        id: String,
658        /// Exec process ID (unique within the container).
659        exec_id: String,
660        /// OCI Process JSON (containerd's ExecProcessRequest spec).
661        process_json: String,
662        /// Allocate a PTY for the exec process.
663        #[serde(default)]
664        tty: bool,
665    },
666
667    /// Signal a pod container's init process (or one exec process).
668    PodSignal {
669        /// Container ID.
670        id: String,
671        /// Exec process to signal instead of init.
672        #[serde(default, skip_serializing_if = "Option::is_none")]
673        exec_id: Option<String>,
674        /// Signal number (SIGKILL = 9, SIGTERM = 15, ...).
675        signal: u32,
676        /// Signal the whole container process group.
677        #[serde(default)]
678        all: bool,
679    },
680
681    /// List PIDs inside a pod container (guest view).
682    PodPids {
683        /// Container ID.
684        id: String,
685    },
686
687    /// Sample a pod container's resource usage (guest view). The agent reads the
688    /// container's process tree from /proc (there is no per-container cgroup); the
689    /// shim maps the reply into containerd's cgroups metrics for CRI stats.
690    PodStats {
691        /// Container ID.
692        id: String,
693    },
694
695    /// Remove a pod container's (or exec process's) guest resources after
696    /// exit: bundle, cgroup, PTY. Exit status was already streamed by
697    /// `PodStart`'s `Exited`.
698    PodDelete {
699        /// Container ID.
700        id: String,
701        /// Exec process to remove instead of the whole container.
702        #[serde(default, skip_serializing_if = "Option::is_none")]
703        exec_id: Option<String>,
704    },
705}
706
707impl AgentRequest {
708    /// A log-safe one-line summary of the request.
709    ///
710    /// This string is written to the machine's console log, which is exposed
711    /// over the logs API — so it must NEVER include credential- or data-bearing
712    /// fields: registry `auth`, `env` (which can carry host-resolved secrets),
713    /// `proxy` (may embed credentials), or `data` (file/stdin bytes). Only the
714    /// variant name plus a non-secret identifier (image) is emitted.
715    ///
716    /// The match is exhaustive with no catch-all on purpose: adding a new
717    /// variant forces a compile error here, so redaction is a deliberate
718    /// decision rather than an accidental leak in some future request type.
719    pub fn log_summary(&self) -> String {
720        match self {
721            AgentRequest::Ping => "Ping".into(),
722            AgentRequest::FsNotify { events } => format!("FsNotify {{ count: {} }}", events.len()),
723            AgentRequest::Pull { image, .. } => format!("Pull {{ image: {image} }}"),
724            AgentRequest::Query { image, .. } => format!("Query {{ image: {image} }}"),
725            AgentRequest::ListImages => "ListImages".into(),
726            AgentRequest::GarbageCollect { .. } => "GarbageCollect".into(),
727            AgentRequest::PrepareOverlay { .. } => "PrepareOverlay".into(),
728            AgentRequest::CleanupOverlay { .. } => "CleanupOverlay".into(),
729            AgentRequest::FormatStorage => "FormatStorage".into(),
730            AgentRequest::StorageStatus => "StorageStatus".into(),
731            AgentRequest::NetworkTest { .. } => "NetworkTest".into(),
732            AgentRequest::Shutdown => "Shutdown".into(),
733            AgentRequest::ExportLayer { .. } => "ExportLayer".into(),
734            AgentRequest::FlattenLayers { lowerdirs, .. } => {
735                format!("FlattenLayers {{ count: {} }}", lowerdirs.len())
736            }
737            AgentRequest::VmExec { .. } => "VmExec".into(),
738            AgentRequest::BranchpointWait { .. } => "BranchpointWait".into(),
739            AgentRequest::BranchpointArm => "BranchpointArm".into(),
740            AgentRequest::BranchpointPark => "BranchpointPark".into(),
741            AgentRequest::BranchpointRelease { .. } => "BranchpointRelease".into(),
742            AgentRequest::BranchpointActivate { .. } => "BranchpointActivate".into(),
743            AgentRequest::BranchpointWaitWorkerReady { .. } => "BranchpointWaitWorkerReady".into(),
744            AgentRequest::Run { image, .. } => format!("Run {{ image: {image} }}"),
745            AgentRequest::Stdin { .. } => "Stdin".into(),
746            AgentRequest::Resize { .. } => "Resize".into(),
747            AgentRequest::FileWrite { .. } => "FileWrite".into(),
748            AgentRequest::FileWriteBegin { .. } => "FileWriteBegin".into(),
749            AgentRequest::FileWriteChunk { .. } => "FileWriteChunk".into(),
750            AgentRequest::FileRead { .. } => "FileRead".into(),
751            AgentRequest::ArchiveDirectory { .. } => "ArchiveDirectory".into(),
752            // Pod requests: spec/process JSON may carry env secrets — emit ids only.
753            AgentRequest::PodCreate { id, .. } => format!("PodCreate {{ id: {id} }}"),
754            AgentRequest::PodStart { id, exec_id } => match exec_id {
755                Some(e) => format!("PodStart {{ id: {id}, exec: {e} }}"),
756                None => format!("PodStart {{ id: {id} }}"),
757            },
758            AgentRequest::PodExec { id, exec_id, .. } => {
759                format!("PodExec {{ id: {id}, exec: {exec_id} }}")
760            }
761            AgentRequest::PodSignal {
762                id, signal, all, ..
763            } => format!("PodSignal {{ id: {id}, signal: {signal}, all: {all} }}"),
764            AgentRequest::PodPids { id } => format!("PodPids {{ id: {id} }}"),
765            AgentRequest::PodStats { id } => format!("PodStats {{ id: {id} }}"),
766            AgentRequest::PodDelete { id, exec_id } => match exec_id {
767                Some(e) => format!("PodDelete {{ id: {id}, exec: {e} }}"),
768                None => format!("PodDelete {{ id: {id} }}"),
769            },
770        }
771    }
772}
773
774/// Agent response types.
775#[derive(Debug, Clone, Serialize, Deserialize)]
776#[serde(tag = "status", rename_all = "snake_case")]
777pub enum AgentResponse {
778    /// Operation completed successfully.
779    Ok {
780        /// Response data (varies by request type).
781        #[serde(default, skip_serializing_if = "Option::is_none")]
782        data: Option<serde_json::Value>,
783    },
784
785    /// Pong response to ping.
786    Pong {
787        /// Protocol version.
788        version: u32,
789        /// Optional agent features that can evolve independently of the base protocol.
790        #[serde(default, skip_serializing_if = "Vec::is_empty")]
791        capabilities: Vec<String>,
792    },
793
794    /// Progress update (for long operations like pull).
795    Progress {
796        /// Human-readable message.
797        message: String,
798        /// Completion percentage (0-100).
799        #[serde(default, skip_serializing_if = "Option::is_none")]
800        percent: Option<u8>,
801        /// Current layer being processed.
802        #[serde(default, skip_serializing_if = "Option::is_none")]
803        layer: Option<String>,
804    },
805
806    /// Operation failed.
807    Error {
808        /// Error message.
809        message: String,
810        /// Error code (for programmatic handling).
811        #[serde(default, skip_serializing_if = "Option::is_none")]
812        code: Option<String>,
813    },
814
815    /// Command execution completed (non-interactive mode).
816    Completed {
817        /// Exit code from the command.
818        exit_code: i32,
819        /// Standard output (may be truncated). `Vec<u8>` preserves binary
820        /// output (image bytes, tarballs, etc.) that would be truncated by
821        /// `String` at the first non-UTF-8 byte. Serialized as base64 JSON
822        /// string — the same format as the streaming `Stdout` variant.
823        #[serde(with = "base64_bytes")]
824        stdout: Vec<u8>,
825        /// Standard error (may be truncated).
826        #[serde(with = "base64_bytes")]
827        stderr: Vec<u8>,
828    },
829
830    /// Command started (interactive mode).
831    /// Indicates the command is running and ready to receive stdin.
832    Started,
833
834    /// Stdout data from a running command (interactive mode).
835    Stdout {
836        /// Output data.
837        #[serde(with = "base64_bytes")]
838        data: Vec<u8>,
839    },
840
841    /// Stderr data from a running command (interactive mode).
842    Stderr {
843        /// Error output data.
844        #[serde(with = "base64_bytes")]
845        data: Vec<u8>,
846    },
847
848    /// Command exited (interactive mode).
849    Exited {
850        /// Exit code from the command.
851        exit_code: i32,
852        /// The container was terminated by the cgroup OOM killer. The shim
853        /// turns this into a TaskOOM event so the CRI reports
854        /// `reason=OOMKilled`. Only ever set on a pod container's init exit.
855        #[serde(default)]
856        oom: bool,
857    },
858
859    /// PIDs inside a pod container (`PodPids` reply).
860    Pids {
861        /// Guest PIDs, container-init first when known.
862        pids: Vec<u32>,
863    },
864
865    /// Resource usage sample for a pod container (`PodStats` reply). Summed over
866    /// the container's process tree read from /proc (no per-container cgroup).
867    Stats {
868        /// Cumulative CPU time of the process tree, in nanoseconds.
869        cpu_usage_ns: u64,
870        /// Resident memory of the process tree, in bytes.
871        memory_bytes: u64,
872    },
873
874    /// Streaming binary-data chunk.
875    ///
876    /// Used by every streaming download direction: the agent sends
877    /// one or more `DataChunk` responses in sequence, with `done: true`
878    /// on the final chunk. Current producers: `ExportLayer` and
879    /// `FileRead`.
880    ///
881    /// Payload size per chunk should stay under
882    /// [`LAYER_CHUNK_SIZE`] so the encoded frame (~1.33× after
883    /// base64) fits inside [`MAX_FRAME_SIZE`] with JSON overhead to
884    /// spare.
885    DataChunk {
886        /// Chunk bytes. Empty allowed on the final frame (common for
887        /// EOF-on-clean-boundary cases).
888        #[serde(with = "base64_bytes")]
889        data: Vec<u8>,
890        /// True on the final chunk of the stream.
891        done: bool,
892    },
893}
894
895// ============================================================================
896// Error Code Constants
897// ============================================================================
898//
899// Standard error codes for AgentResponse::Error. Using constants ensures
900// consistency across the codebase and makes error handling more reliable.
901
902/// Error codes for agent responses.
903pub mod error_codes {
904    /// Request payload was invalid or malformed.
905    pub const INVALID_REQUEST: &str = "INVALID_REQUEST";
906    /// Requested resource was not found.
907    pub const NOT_FOUND: &str = "NOT_FOUND";
908    /// Internal error during operation.
909    pub const INTERNAL_ERROR: &str = "INTERNAL_ERROR";
910    /// Image pull operation failed.
911    pub const PULL_FAILED: &str = "PULL_FAILED";
912    /// Image query operation failed.
913    pub const QUERY_FAILED: &str = "QUERY_FAILED";
914    /// Command execution failed.
915    pub const RUN_FAILED: &str = "RUN_FAILED";
916    /// Command execution failed in container.
917    pub const EXEC_FAILED: &str = "EXEC_FAILED";
918    /// Process spawn failed.
919    pub const SPAWN_FAILED: &str = "SPAWN_FAILED";
920    /// Mount operation failed.
921    pub const MOUNT_FAILED: &str = "MOUNT_FAILED";
922    /// File I/O operation failed.
923    pub const FILE_IO_FAILED: &str = "FILE_IO_FAILED";
924    /// Overlay filesystem operation failed.
925    pub const OVERLAY_FAILED: &str = "OVERLAY_FAILED";
926    /// Cleanup operation failed.
927    pub const CLEANUP_FAILED: &str = "CLEANUP_FAILED";
928    /// Storage format operation failed.
929    pub const FORMAT_FAILED: &str = "FORMAT_FAILED";
930    /// Storage status query failed.
931    pub const STATUS_FAILED: &str = "STATUS_FAILED";
932    /// List operation failed.
933    pub const LIST_FAILED: &str = "LIST_FAILED";
934    /// Garbage collection failed.
935    pub const GC_FAILED: &str = "GC_FAILED";
936    /// Container creation failed.
937    pub const CREATE_FAILED: &str = "CREATE_FAILED";
938    /// Container start failed.
939    pub const START_FAILED: &str = "START_FAILED";
940    /// Container stop failed.
941    pub const STOP_FAILED: &str = "STOP_FAILED";
942    /// Container delete failed.
943    pub const DELETE_FAILED: &str = "DELETE_FAILED";
944    /// Export operation failed.
945    pub const EXPORT_FAILED: &str = "EXPORT_FAILED";
946    /// Serialization error.
947    pub const SERIALIZATION_ERROR: &str = "SERIALIZATION_ERROR";
948    /// Message size exceeds maximum.
949    pub const MESSAGE_TOO_LARGE: &str = "MESSAGE_TOO_LARGE";
950    /// Process wait operation failed.
951    pub const WAIT_FAILED: &str = "WAIT_FAILED";
952}
953
954impl AgentResponse {
955    /// Create an error response with the given message and code.
956    ///
957    /// # Example
958    ///
959    /// ```
960    /// use smolvm_protocol::{AgentResponse, error_codes};
961    ///
962    /// let response = AgentResponse::error("image not found", error_codes::NOT_FOUND);
963    /// ```
964    pub fn error(message: impl Into<String>, code: &str) -> Self {
965        AgentResponse::Error {
966            message: message.into(),
967            code: Some(code.to_string()),
968        }
969    }
970
971    /// Create an error response from a Result's error, with the given code.
972    ///
973    /// # Example
974    ///
975    /// ```ignore
976    /// let response = some_operation()
977    ///     .map(|data| AgentResponse::ok_with_data(data))
978    ///     .unwrap_or_else(|e| AgentResponse::from_err(e, error_codes::PULL_FAILED));
979    /// ```
980    pub fn from_err<E: std::fmt::Display>(err: E, code: &str) -> Self {
981        AgentResponse::Error {
982            message: err.to_string(),
983            code: Some(code.to_string()),
984        }
985    }
986
987    /// Create an Ok response with optional JSON data.
988    pub fn ok(data: Option<serde_json::Value>) -> Self {
989        AgentResponse::Ok { data }
990    }
991
992    /// Create an Ok response with JSON-serializable data.
993    ///
994    /// Returns an error response if serialization fails.
995    pub fn ok_with_data<T: serde::Serialize>(data: T) -> Self {
996        match serde_json::to_value(data) {
997            Ok(value) => AgentResponse::Ok { data: Some(value) },
998            Err(e) => AgentResponse::error(
999                format!("failed to serialize response: {}", e),
1000                error_codes::SERIALIZATION_ERROR,
1001            ),
1002        }
1003    }
1004
1005    /// Convert a Result into an AgentResponse.
1006    ///
1007    /// On success, serializes the value to JSON. On error, creates an error response.
1008    ///
1009    /// # Example
1010    ///
1011    /// ```ignore
1012    /// let response = AgentResponse::from_result(
1013    ///     storage::pull_image(image),
1014    ///     error_codes::PULL_FAILED,
1015    /// );
1016    /// ```
1017    pub fn from_result<T, E>(result: Result<T, E>, error_code: &str) -> Self
1018    where
1019        T: serde::Serialize,
1020        E: std::fmt::Display,
1021    {
1022        match result {
1023            Ok(data) => Self::ok_with_data(data),
1024            Err(e) => Self::from_err(e, error_code),
1025        }
1026    }
1027}
1028
1029/// Image information returned by Query/ListImages.
1030#[derive(Debug, Clone, Serialize, Deserialize)]
1031pub struct ImageInfo {
1032    /// Image reference.
1033    pub reference: String,
1034    /// Image digest (sha256:...).
1035    pub digest: String,
1036    /// Image size in bytes.
1037    pub size: u64,
1038    /// Creation timestamp (ISO 8601).
1039    pub created: Option<String>,
1040    /// Platform architecture.
1041    pub architecture: String,
1042    /// Platform OS.
1043    pub os: String,
1044    /// Number of layers.
1045    pub layer_count: usize,
1046    /// Layer digests in order.
1047    pub layers: Vec<String>,
1048    /// Image entrypoint (from OCI config).
1049    #[serde(default)]
1050    pub entrypoint: Vec<String>,
1051    /// Image default command (from OCI config).
1052    #[serde(default)]
1053    pub cmd: Vec<String>,
1054    /// Image environment variables (from OCI config).
1055    #[serde(default)]
1056    pub env: Vec<String>,
1057    /// Image working directory (from OCI config).
1058    #[serde(default)]
1059    pub workdir: Option<String>,
1060    /// Image default user (from OCI config).
1061    #[serde(default)]
1062    pub user: Option<String>,
1063}
1064
1065/// Overlay preparation result.
1066#[derive(Debug, Clone, Serialize, Deserialize)]
1067pub struct OverlayInfo {
1068    /// Path to the merged overlay rootfs.
1069    pub rootfs_path: String,
1070    /// Path to the upper (writable) directory.
1071    pub upper_path: String,
1072    /// Path to the work directory.
1073    pub work_path: String,
1074}
1075
1076/// Storage status information.
1077#[derive(Debug, Clone, Serialize, Deserialize)]
1078pub struct StorageStatus {
1079    /// Whether the storage is formatted and ready.
1080    pub ready: bool,
1081    /// Total size in bytes.
1082    pub total_bytes: u64,
1083    /// Used size in bytes.
1084    pub used_bytes: u64,
1085    /// Number of cached layers.
1086    pub layer_count: usize,
1087    /// Number of cached images.
1088    pub image_count: usize,
1089}
1090
1091/// Registry authentication credentials for pulling images.
1092///
1093/// `Debug` is hand-written to redact the password: this value is carried inside
1094/// `AgentRequest::Pull`, and any `{:?}` of that request (e.g. a tracing span)
1095/// would otherwise serialize the token verbatim into the machine's console log,
1096/// which is exposed over the logs API.
1097#[derive(Clone, Serialize, Deserialize)]
1098pub struct RegistryAuth {
1099    /// Username for authentication.
1100    pub username: String,
1101    /// Password or token for authentication.
1102    pub password: String,
1103}
1104
1105impl std::fmt::Debug for RegistryAuth {
1106    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1107        f.debug_struct("RegistryAuth")
1108            .field("username", &self.username)
1109            .field("password", &"***")
1110            .finish()
1111    }
1112}
1113
1114// ============================================================================
1115// Workload VM Protocol (Command Execution)
1116// ============================================================================
1117
1118/// Messages from host to workload VM.
1119#[derive(Debug, Clone, Serialize, Deserialize)]
1120#[serde(tag = "type", rename_all = "snake_case")]
1121pub enum HostMessage {
1122    /// Authentication request.
1123    Auth {
1124        /// Authentication token (base64).
1125        token: String,
1126        /// Protocol version.
1127        protocol_version: u32,
1128    },
1129
1130    /// Run a command.
1131    Run {
1132        /// Request ID for correlating responses.
1133        request_id: u64,
1134        /// Command and arguments.
1135        command: Vec<String>,
1136        /// Environment variables.
1137        env: Vec<(String, String)>,
1138        /// Working directory.
1139        workdir: Option<String>,
1140    },
1141
1142    /// Execute a command in running VM.
1143    Exec {
1144        /// Request ID.
1145        request_id: u64,
1146        /// Command and arguments.
1147        command: Vec<String>,
1148        /// Allocate a TTY.
1149        tty: bool,
1150    },
1151
1152    /// Send a signal to a running command.
1153    Signal {
1154        /// Request ID of the command.
1155        request_id: u64,
1156        /// Signal number.
1157        signal: i32,
1158    },
1159
1160    /// Request graceful shutdown.
1161    Stop {
1162        /// Timeout in milliseconds.
1163        timeout_ms: u64,
1164    },
1165}
1166
1167/// Messages from workload VM to host.
1168#[derive(Debug, Clone, Serialize, Deserialize)]
1169#[serde(tag = "type", rename_all = "snake_case")]
1170pub enum GuestMessage {
1171    /// Authentication successful.
1172    AuthOk,
1173
1174    /// Authentication failed.
1175    AuthFailed,
1176
1177    /// VM is ready to receive commands.
1178    Ready,
1179
1180    /// Command started.
1181    Started {
1182        /// Request ID.
1183        request_id: u64,
1184    },
1185
1186    /// Stdout data from command.
1187    Stdout {
1188        /// Request ID.
1189        request_id: u64,
1190        /// Output data.
1191        #[serde(with = "base64_bytes")]
1192        data: Vec<u8>,
1193        /// Whether output was truncated.
1194        truncated: bool,
1195    },
1196
1197    /// Stderr data from command.
1198    Stderr {
1199        /// Request ID.
1200        request_id: u64,
1201        /// Output data.
1202        #[serde(with = "base64_bytes")]
1203        data: Vec<u8>,
1204        /// Whether output was truncated.
1205        truncated: bool,
1206    },
1207
1208    /// Command exited.
1209    Exit {
1210        /// Request ID.
1211        request_id: u64,
1212        /// Exit code.
1213        code: i32,
1214        /// Exit reason.
1215        reason: String,
1216    },
1217
1218    /// Error occurred.
1219    Error {
1220        /// Request ID (if applicable).
1221        request_id: Option<u64>,
1222        /// Error message.
1223        message: String,
1224    },
1225}
1226
1227// ============================================================================
1228// Wire Format Helpers
1229// ============================================================================
1230
1231/// Envelope that wraps any message with an optional trace ID for correlation.
1232///
1233/// On the wire, the trace_id is flattened into the JSON alongside the message
1234/// fields: `{"trace_id":"abc123","method":"ping"}`.
1235#[derive(Debug, Clone, Serialize, Deserialize)]
1236pub struct Envelope<T> {
1237    /// Trace ID for correlating host API requests to agent operations.
1238    #[serde(skip_serializing_if = "Option::is_none", default)]
1239    pub trace_id: Option<String>,
1240    /// The wrapped message.
1241    #[serde(flatten)]
1242    pub body: T,
1243}
1244
1245impl<T> Envelope<T> {
1246    /// Create an envelope with no trace ID.
1247    pub fn new(body: T) -> Self {
1248        Self {
1249            trace_id: None,
1250            body,
1251        }
1252    }
1253
1254    /// Create an envelope with an optional trace ID.
1255    pub fn with_trace_id(body: T, trace_id: Option<String>) -> Self {
1256        Self { trace_id, body }
1257    }
1258}
1259
1260/// Encode a message to wire format (length-prefixed JSON).
1261pub fn encode_message<T: Serialize>(msg: &T) -> Result<Vec<u8>, serde_json::Error> {
1262    let json = serde_json::to_vec(msg)?;
1263    let len = json.len() as u32;
1264
1265    let mut buf = Vec::with_capacity(4 + json.len());
1266    buf.extend_from_slice(&len.to_be_bytes());
1267    buf.extend_from_slice(&json);
1268
1269    Ok(buf)
1270}
1271
1272/// Decode a message from wire format.
1273pub fn decode_message<T: for<'de> Deserialize<'de>>(data: &[u8]) -> Result<T, DecodeError> {
1274    if data.len() < 4 {
1275        return Err(DecodeError::TooShort);
1276    }
1277
1278    let len = u32::from_be_bytes([data[0], data[1], data[2], data[3]]) as usize;
1279
1280    if len > MAX_FRAME_SIZE as usize {
1281        return Err(DecodeError::TooLarge(len));
1282    }
1283
1284    if data.len() < 4 + len {
1285        return Err(DecodeError::Incomplete {
1286            expected: len,
1287            got: data.len() - 4,
1288        });
1289    }
1290
1291    serde_json::from_slice(&data[4..4 + len]).map_err(DecodeError::Json)
1292}
1293
1294/// Error decoding a wire message.
1295#[derive(Debug)]
1296pub enum DecodeError {
1297    /// Data too short to contain length header.
1298    TooShort,
1299    /// Frame size exceeds maximum.
1300    TooLarge(usize),
1301    /// Incomplete frame.
1302    Incomplete {
1303        /// Expected length.
1304        expected: usize,
1305        /// Actual length.
1306        got: usize,
1307    },
1308    /// JSON parse error.
1309    Json(serde_json::Error),
1310}
1311
1312impl std::fmt::Display for DecodeError {
1313    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1314        match self {
1315            DecodeError::TooShort => write!(f, "data too short for length header"),
1316            DecodeError::TooLarge(size) => write!(f, "frame too large: {} bytes", size),
1317            DecodeError::Incomplete { expected, got } => {
1318                write!(
1319                    f,
1320                    "incomplete frame: expected {} bytes, got {}",
1321                    expected, got
1322                )
1323            }
1324            DecodeError::Json(e) => write!(f, "JSON decode error: {}", e),
1325        }
1326    }
1327}
1328
1329impl std::error::Error for DecodeError {}
1330
1331#[cfg(test)]
1332mod tests {
1333    use super::*;
1334
1335    #[test]
1336    fn file_write_without_owner_fields_still_parses() {
1337        // Requests from clients predating uid/gid must keep deserializing.
1338        let old = r#"{"method":"file_write","path":"/x","data":"aGk=","mode":420}"#;
1339        let req: AgentRequest = serde_json::from_str(old).unwrap();
1340        match req {
1341            AgentRequest::FileWrite { mode, uid, gid, .. } => {
1342                assert_eq!(mode, Some(420));
1343                assert_eq!(uid, None);
1344                assert_eq!(gid, None);
1345            }
1346            other => panic!("unexpected: {other:?}"),
1347        }
1348    }
1349
1350    #[test]
1351    fn test_encode_decode_roundtrip() {
1352        let req = AgentRequest::Pull {
1353            image: "alpine:latest".to_string(),
1354            oci_platform: Some("linux/arm64".to_string()),
1355            auth: None,
1356            proxy: None,
1357            no_proxy: None,
1358        };
1359
1360        let encoded = encode_message(&req).unwrap();
1361        let decoded: AgentRequest = decode_message(&encoded).unwrap();
1362
1363        let AgentRequest::Pull {
1364            image,
1365            oci_platform,
1366            auth,
1367            proxy,
1368            no_proxy,
1369        } = decoded
1370        else {
1371            panic!("expected Pull variant, got {:?}", decoded);
1372        };
1373        assert_eq!(image, "alpine:latest");
1374        assert_eq!(oci_platform, Some("linux/arm64".to_string()));
1375        assert!(auth.is_none());
1376        assert!(proxy.is_none());
1377        assert!(no_proxy.is_none());
1378    }
1379
1380    #[test]
1381    fn test_encode_decode_with_auth() {
1382        let req = AgentRequest::Pull {
1383            image: "ghcr.io/owner/repo:latest".to_string(),
1384            oci_platform: None,
1385            auth: Some(RegistryAuth {
1386                username: "testuser".to_string(),
1387                password: "testpass".to_string(),
1388            }),
1389            proxy: None,
1390            no_proxy: None,
1391        };
1392
1393        let encoded = encode_message(&req).unwrap();
1394        let decoded: AgentRequest = decode_message(&encoded).unwrap();
1395
1396        let AgentRequest::Pull {
1397            image,
1398            oci_platform,
1399            auth,
1400            proxy: _,
1401            no_proxy: _,
1402        } = decoded
1403        else {
1404            panic!("expected Pull variant, got {:?}", decoded);
1405        };
1406        assert_eq!(image, "ghcr.io/owner/repo:latest");
1407        assert!(oci_platform.is_none());
1408        let auth = auth.expect("auth should be Some");
1409        assert_eq!(auth.username, "testuser");
1410        assert_eq!(auth.password, "testpass");
1411    }
1412
1413    #[test]
1414    fn test_encode_decode_with_proxy() {
1415        let req = AgentRequest::Pull {
1416            image: "alpine:latest".to_string(),
1417            oci_platform: None,
1418            auth: None,
1419            proxy: Some("http://192.168.127.254:3128".to_string()),
1420            no_proxy: Some("127.0.0.1,localhost,.internal".to_string()),
1421        };
1422
1423        let encoded = encode_message(&req).unwrap();
1424        let decoded: AgentRequest = decode_message(&encoded).unwrap();
1425
1426        let AgentRequest::Pull {
1427            proxy, no_proxy, ..
1428        } = decoded
1429        else {
1430            panic!("expected Pull variant, got {:?}", decoded);
1431        };
1432        assert_eq!(proxy.as_deref(), Some("http://192.168.127.254:3128"));
1433        assert_eq!(no_proxy.as_deref(), Some("127.0.0.1,localhost,.internal"));
1434    }
1435
1436    #[test]
1437    fn test_decode_too_short() {
1438        let data = [0u8; 2];
1439        let result: Result<AgentRequest, _> = decode_message(&data);
1440        assert!(matches!(result, Err(DecodeError::TooShort)));
1441    }
1442
1443    #[test]
1444    fn test_decode_incomplete() {
1445        let mut data = vec![0, 0, 0, 100]; // claims 100 bytes
1446        data.extend_from_slice(b"{}"); // only 2 bytes of payload
1447        let result: Result<AgentRequest, _> = decode_message(&data);
1448        assert!(matches!(result, Err(DecodeError::Incomplete { .. })));
1449    }
1450
1451    #[test]
1452    fn test_agent_request_serialization() {
1453        let req = AgentRequest::Ping;
1454        let json = serde_json::to_string(&req).unwrap();
1455        assert!(json.contains("ping"));
1456
1457        let req = AgentRequest::PrepareOverlay {
1458            image: "ubuntu:22.04".to_string(),
1459            workload_id: "wl-123".to_string(),
1460        };
1461        let json = serde_json::to_string(&req).unwrap();
1462        assert!(json.contains("prepare_overlay"));
1463    }
1464
1465    #[test]
1466    fn test_agent_response_serialization() {
1467        let resp = AgentResponse::Pong {
1468            version: PROTOCOL_VERSION,
1469            capabilities: vec![forkpoint::TYPED_BRANCHPOINT_CAPABILITY.to_string()],
1470        };
1471        let json = serde_json::to_string(&resp).unwrap();
1472        assert!(json.contains("pong"));
1473        assert!(json.contains(forkpoint::TYPED_BRANCHPOINT_CAPABILITY));
1474
1475        let legacy: AgentResponse =
1476            serde_json::from_str(r#"{"status":"pong","version":1}"#).unwrap();
1477        assert!(matches!(
1478            legacy,
1479            AgentResponse::Pong {
1480                version: PROTOCOL_VERSION,
1481                capabilities
1482            } if capabilities.is_empty()
1483        ));
1484
1485        let resp = AgentResponse::Progress {
1486            message: "Pulling layer 1/3".to_string(),
1487            percent: Some(33),
1488            layer: Some("sha256:abc123".to_string()),
1489        };
1490        let json = serde_json::to_string(&resp).unwrap();
1491        assert!(json.contains("progress"));
1492    }
1493
1494    #[test]
1495    fn file_write_begin_roundtrips() {
1496        let req = AgentRequest::FileWriteBegin {
1497            path: "/tmp/target".into(),
1498            mode: Some(0o600),
1499            uid: Some(1000),
1500            gid: Some(1000),
1501            total_size: 123_456_789,
1502        };
1503        let bytes = encode_message(&req).unwrap();
1504        let back: AgentRequest = decode_message(&bytes).unwrap();
1505        match back {
1506            AgentRequest::FileWriteBegin {
1507                path,
1508                mode,
1509                uid,
1510                gid,
1511                total_size,
1512            } => {
1513                assert_eq!(path, "/tmp/target");
1514                assert_eq!(mode, Some(0o600));
1515                assert_eq!(uid, Some(1000));
1516                assert_eq!(gid, Some(1000));
1517                assert_eq!(total_size, 123_456_789);
1518            }
1519            _ => panic!("wrong variant"),
1520        }
1521    }
1522
1523    #[test]
1524    fn file_write_chunk_roundtrips_binary_data() {
1525        // Binary data (bytes outside UTF-8) must survive the base64
1526        // trip intact. If the encoding ever silently lossifies, this
1527        // fires.
1528        let payload: Vec<u8> = (0u8..=255).collect();
1529        let req = AgentRequest::FileWriteChunk {
1530            data: payload.clone(),
1531            done: true,
1532        };
1533        let bytes = encode_message(&req).unwrap();
1534        let back: AgentRequest = decode_message(&bytes).unwrap();
1535        match back {
1536            AgentRequest::FileWriteChunk { data, done } => {
1537                assert_eq!(data, payload);
1538                assert!(done);
1539            }
1540            _ => panic!("wrong variant"),
1541        }
1542    }
1543
1544    #[test]
1545    fn file_write_size_constants_are_frame_safe() {
1546        // Sanity: a single streaming chunk at FILE_WRITE_CHUNK_SIZE
1547        // must fit inside MAX_FRAME_SIZE after base64 (+ ~33%) and
1548        // JSON overhead. If anyone bumps CHUNK_SIZE past the limit,
1549        // this test fires before production does.
1550        let chunk_bytes = FILE_WRITE_CHUNK_SIZE as u64;
1551        let base64_bytes = chunk_bytes.div_ceil(3) * 4; // ceil(n/3)*4
1552        let json_overhead = 256u64; // method tag, done bool, quotes
1553        let total = base64_bytes + json_overhead;
1554        assert!(
1555            total < MAX_FRAME_SIZE as u64,
1556            "FILE_WRITE_CHUNK_SIZE of {} bytes would produce a frame \
1557             of ~{} bytes which exceeds MAX_FRAME_SIZE of {}",
1558            chunk_bytes,
1559            total,
1560            MAX_FRAME_SIZE
1561        );
1562    }
1563
1564    #[test]
1565    fn test_ports_constants() {
1566        assert_eq!(ports::WORKLOAD_CONTROL, 5000);
1567        assert_eq!(ports::WORKLOAD_LOGS, 5001);
1568        assert_eq!(ports::AGENT_CONTROL, 6000);
1569        assert_eq!(ports::SSH_AGENT, 6001);
1570    }
1571
1572    #[test]
1573    fn test_cid_constants() {
1574        assert_eq!(cid::HOST, 2);
1575        assert_eq!(cid::GUEST, 3);
1576    }
1577
1578    #[test]
1579    fn test_envelope_serialization_with_trace_id() {
1580        let req = AgentRequest::Ping;
1581        let envelope = Envelope::with_trace_id(&req, Some("abc123".to_string()));
1582        let json = serde_json::to_string(&envelope).unwrap();
1583
1584        // trace_id should be flattened alongside the method tag
1585        assert!(json.contains("\"trace_id\":\"abc123\""));
1586        assert!(json.contains("\"method\":\"ping\""));
1587
1588        // Deserialize back — Envelope<AgentRequest> with flatten
1589        let parsed: Envelope<AgentRequest> = serde_json::from_str(&json).unwrap();
1590        assert_eq!(parsed.trace_id.as_deref(), Some("abc123"));
1591        assert!(matches!(parsed.body, AgentRequest::Ping));
1592    }
1593
1594    #[test]
1595    fn test_envelope_without_trace_id() {
1596        let req = AgentRequest::Ping;
1597        let envelope = Envelope::new(&req);
1598        let json = serde_json::to_string(&envelope).unwrap();
1599
1600        // No trace_id field (skip_serializing_if = None)
1601        assert!(!json.contains("trace_id"));
1602        assert!(json.contains("\"method\":\"ping\""));
1603    }
1604
1605    #[test]
1606    fn test_envelope_backward_compat_bare_request() {
1607        // A bare AgentRequest (no Envelope) should fail to parse as Envelope
1608        // but succeed as bare AgentRequest — this is the agent's fallback path
1609        let bare_json = r#"{"method":"ping"}"#;
1610
1611        // Envelope parse should fail (no body field to flatten into)
1612        // Actually with flatten, this may work — let's verify
1613        let envelope_result = serde_json::from_str::<Envelope<AgentRequest>>(bare_json);
1614        let bare_result = serde_json::from_str::<AgentRequest>(bare_json);
1615
1616        // At least one must succeed for backward compat
1617        assert!(
1618            envelope_result.is_ok() || bare_result.is_ok(),
1619            "Neither Envelope nor bare parse succeeded"
1620        );
1621
1622        // Bare parse must always work
1623        assert!(bare_result.is_ok());
1624        assert!(matches!(bare_result.unwrap(), AgentRequest::Ping));
1625
1626        // If Envelope works, trace_id should be None
1627        if let Ok(env) = envelope_result {
1628            assert!(env.trace_id.is_none());
1629        }
1630    }
1631}