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