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