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