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