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