Skip to main content

smolvm_protocol/
lib.rs

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