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