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