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