mod cancellation;
mod machine_exec;
mod sandbox_stream;
mod shutdown;
mod storage_check;
mod storage_health;
mod transport;
mod wire;
#[cfg(test)]
mod admission_tests;
pub use self::machine_exec::{ExecSessionInput, ExecSessionOutput};
pub use self::sandbox_stream::SandboxStream;
use self::sandbox_stream::StreamKind;
pub use self::storage_health::StorageHealthStream;
use self::transport::{AgentTransport, BLOCKING_RPC_TIMEOUT};
use crate::error::{EngineError, Result};
use arcbox_connect::sandbox_v1::{
AttachExecutionRequest, BuildTemplateRequest, CheckpointRequest, CheckpointResponse,
CreateSandboxRequest, CreateSandboxResponse, DeleteSnapshotRequest, DeleteTemplateRequest,
Execution, ExecutionEvent, FileChunk, FileStat, GetStdinStatusRequest, GetTemplateRequest,
InspectSandboxRequest, ListDirRequest, ListDirResponse, ListExecutionsRequest,
ListExecutionsResponse, ListSandboxesRequest, ListSandboxesResponse, ListSnapshotsRequest,
ListSnapshotsResponse, ListTemplatesRequest, ListTemplatesResponse, MakeDirRequest,
MoveEntryRequest, PauseSandboxRequest, PublishTemplateRequest, ReadFileRequest,
RemoveEntryRequest, RemoveSandboxRequest, ResizeExecutionTtyRequest, RestoreRequest,
RestoreResponse, SandboxEvent, SandboxEventsRequest, SandboxInfo, SetLifecycleRequest,
SignalExecutionRequest, StartExecutionRequest, StatFileRequest, StdinStatus,
StopSandboxRequest, Template, WaitExecutionRequest, WaitForPortRequest, WatchDirRequest,
WatchDirResponse, WriteFileOpen, WriteStdinRequest, execution_event,
};
use arcbox_connect::v1::{
AgentPingRequest as PingRequest, AgentPingResponse as PingResponse, ContainerFsPathsRequest,
ContainerFsPathsResponse, DiskTrimRequest, DiskTrimResponse, EnsureMachineExportRequest,
EnsureMachineExportResponse, EnsureNfsExportRequest, EnsureNfsExportResponse,
ImageFsPathsRequest, ImageFsPathsResponse, KubernetesDeleteRequest, KubernetesDeleteResponse,
KubernetesKubeconfigRequest, KubernetesKubeconfigResponse, KubernetesLoadBalancersRequest,
KubernetesLoadBalancersResponse, KubernetesStartRequest, KubernetesStartResponse,
KubernetesStatusRequest, KubernetesStatusResponse, KubernetesStopRequest,
KubernetesStopResponse, MachineStats, MemoryPressureEvent, MmapReadFileRequest,
MmapReadFileResponse, ReadinessEvent, RuntimeEnsureRequest, RuntimeEnsureResponse,
RuntimeStatusRequest, RuntimeStatusResponse, SandboxCleanupResponse, SandboxCleanupTicket,
SandboxPortForwardRemoveRequest, SandboxPortForwardRequest, SandboxPortForwardResponse,
SandboxResumeCommand, SandboxResumeResponse, SystemInfo, WatchMemoryPressureRequest,
WatchReadinessRequest, WatchSandboxCleanupRequest, WatchStatsRequest,
};
use arcbox_constants::ports::AGENT_PORT;
use arcbox_constants::wire::MessageType;
use arcbox_transport::Transport;
use arcbox_transport::vsock::BlockingVsockTransport;
use arcbox_transport::vsock::{VsockAddr, VsockTransport};
use buffa::Message;
use bytes::Bytes;
use std::time::{Duration, Instant};
use tokio::sync::mpsc;
#[derive(Debug)]
pub enum WriteFileChunk {
Data(Vec<u8>),
Abort,
}
pub struct AgentClient {
cid: u32,
transport: AgentTransport,
connected: bool,
protocol_admitted: bool,
}
struct PingHandshake<'a> {
client: &'a mut AgentClient,
completed: bool,
}
impl<'a> PingHandshake<'a> {
fn new(client: &'a mut AgentClient) -> Self {
client.protocol_admitted = false;
Self {
client,
completed: false,
}
}
fn finish(mut self, kind: u32, payload: &[u8]) -> Result<PingResponse> {
let response = self.client.decode_ping_response(kind, payload)?;
self.completed = true;
Ok(response)
}
}
impl Drop for PingHandshake<'_> {
fn drop(&mut self) {
if !self.completed {
self.client.transport.close();
self.client.connected = false;
}
}
}
impl AgentClient {
#[must_use]
pub const fn new(cid: u32) -> Self {
let addr = VsockAddr::new(cid, AGENT_PORT);
Self {
cid,
transport: AgentTransport::Async(VsockTransport::new(addr)),
connected: false,
protocol_admitted: false,
}
}
pub fn from_fd_blocking(cid: u32, fd: std::os::unix::io::RawFd) -> Result<Self> {
let transport = unsafe { BlockingVsockTransport::from_raw_fd(fd) }
.map_err(|e| EngineError::Machine(format!("invalid vsock fd: {e}")))?;
Ok(Self {
cid,
transport: AgentTransport::Blocking(transport),
connected: true,
protocol_admitted: false,
})
}
#[cfg(target_os = "macos")]
pub fn from_fd_async(cid: u32, fd: std::os::unix::io::RawFd) -> Result<Self> {
let addr = VsockAddr::new(cid, AGENT_PORT);
let transport = VsockTransport::from_raw_fd(fd, addr)
.map_err(|e| EngineError::Machine(format!("invalid vsock fd: {e}")))?;
Ok(Self {
cid,
transport: AgentTransport::Async(transport),
connected: true,
protocol_admitted: false,
})
}
#[must_use]
pub const fn cid(&self) -> u32 {
self.cid
}
pub fn build_message(msg_type: MessageType, trace_id: &str, payload: &[u8]) -> Bytes {
wire::build_message(msg_type, trace_id, payload)
}
pub async fn connect(&mut self) -> Result<()> {
if self.connected {
return Ok(());
}
self.protocol_admitted = false;
match &mut self.transport {
AgentTransport::Async(t) => {
t.connect().await.map_err(|source| EngineError::Transport {
context: "failed to connect to agent",
source,
})?;
}
AgentTransport::Blocking(_) => {
return Err(EngineError::Machine(
"blocking agent connection is closed; obtain a new VM socket".into(),
));
}
}
self.connected = true;
tracing::debug!(cid = self.cid, "connected to agent");
Ok(())
}
pub async fn disconnect(&mut self) -> Result<()> {
self.protocol_admitted = false;
self.transport.close();
self.connected = false;
Ok(())
}
async fn rpc_call(&mut self, msg_type: MessageType, payload: &[u8]) -> Result<(u32, Vec<u8>)> {
let trace_id = crate::trace::current_trace_id();
self.rpc_call_traced(msg_type, &trace_id, payload).await
}
async fn rpc_call_traced(
&mut self,
msg_type: MessageType,
trace_id: &str,
payload: &[u8],
) -> Result<(u32, Vec<u8>)> {
self.require_agent_protocol().await?;
self.rpc_exchange_traced(msg_type, trace_id, payload).await
}
async fn rpc_exchange_traced(
&mut self,
msg_type: MessageType,
trace_id: &str,
payload: &[u8],
) -> Result<(u32, Vec<u8>)> {
if !self.connected {
self.connect().await?;
}
let buf = wire::build_message(msg_type, trace_id, payload);
let response = match &mut self.transport {
AgentTransport::Async(t) => {
t.send(buf)
.await
.map_err(|e| EngineError::Machine(format!("{msg_type:?} send failed: {e}")))?;
t.recv().await.map_err(|e| {
EngineError::Machine(format!("{msg_type:?} response receive failed: {e}"))
})?
}
AgentTransport::Blocking(t) => {
tokio::task::block_in_place(|| {
let deadline = Instant::now() + BLOCKING_RPC_TIMEOUT;
t.send(&buf, deadline).map_err(|e| {
EngineError::Machine(format!("{msg_type:?} send failed: {e}"))
})?;
t.recv(deadline).map_err(|e| {
EngineError::Machine(format!("{msg_type:?} response receive failed: {e}"))
})
})?
}
};
let (resp_type, _resp_trace, payload) = wire::parse_response(&response)?;
if resp_type == MessageType::Error as u32 {
let (code, message) = wire::parse_error_response(&payload)?;
return Err(EngineError::Agent { code, message });
}
Ok((resp_type, payload))
}
fn rpc_call_blocking(
&mut self,
msg_type: MessageType,
payload: &[u8],
) -> Result<(u32, Vec<u8>)> {
self.require_agent_protocol_blocking()?;
self.rpc_exchange_blocking(msg_type, payload)
}
fn rpc_exchange_blocking(
&mut self,
msg_type: MessageType,
payload: &[u8],
) -> Result<(u32, Vec<u8>)> {
let trace_id = "";
let buf = wire::build_message(msg_type, trace_id, payload);
let response = match &mut self.transport {
AgentTransport::Blocking(t) => {
let deadline = Instant::now() + BLOCKING_RPC_TIMEOUT;
t.send(&buf, deadline)
.map_err(|e| EngineError::Machine(format!("{msg_type:?} send failed: {e}")))?;
t.recv(deadline).map_err(|e| {
EngineError::Machine(format!("{msg_type:?} response receive failed: {e}"))
})?
}
AgentTransport::Async(_) => {
return Err(EngineError::Machine(
"rpc_call_blocking called on async transport".into(),
));
}
};
let (resp_type, _resp_trace, payload) = wire::parse_response(&response)?;
if resp_type == MessageType::Error as u32 {
let (code, message) = wire::parse_error_response(&payload)?;
return Err(EngineError::Agent { code, message });
}
Ok((resp_type, payload))
}
fn decode_response<T: Message>(payload: &[u8]) -> Result<T> {
T::decode_from_slice(payload)
.map_err(|e| EngineError::Machine(format!("failed to decode response: {e}")))
}
fn expect_response_type(resp_type: u32, expected: MessageType) -> Result<()> {
if resp_type == expected as u32 {
Ok(())
} else {
Err(EngineError::Machine(format!(
"unexpected response type: 0x{resp_type:04x}"
)))
}
}
fn expect_ack_response_type(resp_type: u32, expected: MessageType) -> Result<()> {
if resp_type == expected as u32 || resp_type == MessageType::Empty as u32 {
Ok(())
} else {
Err(EngineError::Machine(format!(
"unexpected response type: 0x{resp_type:04x}"
)))
}
}
async fn unary_rpc<T: Message>(
&mut self,
request_type: MessageType,
payload: &[u8],
response_type: MessageType,
) -> Result<T> {
let (resp_type, resp_payload) = self.rpc_call(request_type, payload).await?;
Self::expect_response_type(resp_type, response_type)?;
Self::decode_response(&resp_payload)
}
fn unary_rpc_blocking<T: Message>(
&mut self,
request_type: MessageType,
payload: &[u8],
response_type: MessageType,
) -> Result<T> {
let (resp_type, resp_payload) = self.rpc_call_blocking(request_type, payload)?;
Self::expect_response_type(resp_type, response_type)?;
Self::decode_response(&resp_payload)
}
pub fn check_agent_protocol(resp: &PingResponse) -> Result<()> {
use arcbox_constants::wire::{AGENT_PROTOCOL_VERSION, MIN_AGENT_PROTOCOL_VERSION};
if resp.protocol_version == 0 {
return Err(EngineError::Machine(format!(
"guest agent did not complete a compatible handshake \
(agent version {}, response {:?}); fix the reported guest boot contract",
resp.version, resp.message,
)));
}
if resp.protocol_version < MIN_AGENT_PROTOCOL_VERSION {
return Err(EngineError::Machine(format!(
"guest agent is incompatible with this daemon: agent protocol {} \
(agent version {}, response {:?}), daemon requires >= {}. The staged agent \
binary is stale — reinstall or update ArcBox so the bundled \
agent is staged again",
resp.protocol_version, resp.version, resp.message, MIN_AGENT_PROTOCOL_VERSION,
)));
}
if resp.protocol_version > AGENT_PROTOCOL_VERSION {
tracing::warn!(
agent_protocol = resp.protocol_version,
host_protocol = AGENT_PROTOCOL_VERSION,
agent_version = %resp.version,
"guest agent speaks a newer protocol than this daemon; \
continuing (protocol evolution is additive)"
);
}
Ok(())
}
async fn require_agent_protocol(&mut self) -> Result<()> {
if self.protocol_admitted {
return Ok(());
}
let response = tokio::time::timeout(BLOCKING_RPC_TIMEOUT, self.ping())
.await
.map_err(|_| EngineError::Machine("agent protocol handshake timed out".to_owned()))??;
if self.protocol_admitted {
Ok(())
} else {
Self::check_agent_protocol(&response)
}
}
fn require_agent_protocol_blocking(&mut self) -> Result<()> {
if self.protocol_admitted {
return Ok(());
}
let response = self.ping_blocking()?;
if self.protocol_admitted {
Ok(())
} else {
Self::check_agent_protocol(&response)
}
}
fn decode_ping_response(&mut self, kind: u32, payload: &[u8]) -> Result<PingResponse> {
Self::expect_response_type(kind, MessageType::PingResponse)?;
let response = Self::decode_response(payload)?;
self.protocol_admitted = Self::check_agent_protocol(&response).is_ok();
Ok(response)
}
pub fn ping_blocking(&mut self) -> Result<PingResponse> {
let handshake = PingHandshake::new(self);
let req = PingRequest {
message: "ping".to_string(),
timestamp_secs: std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_or(0, |d| i64::try_from(d.as_secs()).unwrap_or(0)),
..Default::default()
};
let payload = req.encode_to_vec();
let (kind, payload) = handshake
.client
.rpc_exchange_blocking(MessageType::PingRequest, &payload)?;
handshake.finish(kind, &payload)
}
pub fn is_blocking(&self) -> bool {
matches!(self.transport, AgentTransport::Blocking(_))
}
pub fn kill_agent_blocking(&mut self) -> Result<()> {
let (resp_type, _payload) = self.rpc_call_blocking(MessageType::KillAgentRequest, &[])?;
if resp_type != MessageType::KillAgentResponse as u32 {
return Err(EngineError::Machine(format!(
"unexpected response type: {resp_type}"
)));
}
Ok(())
}
pub async fn ping(&mut self) -> Result<PingResponse> {
let handshake = PingHandshake::new(self);
let req = PingRequest {
message: "ping".to_string(),
timestamp_secs: std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_or(0, |d| i64::try_from(d.as_secs()).unwrap_or(0)),
..Default::default()
};
let payload = req.encode_to_vec();
let trace_id = crate::trace::current_trace_id();
let (kind, payload) = handshake
.client
.rpc_exchange_traced(MessageType::PingRequest, &trace_id, &payload)
.await?;
handshake.finish(kind, &payload)
}
pub fn get_system_info_blocking(&mut self) -> Result<SystemInfo> {
self.unary_rpc_blocking(
MessageType::GetSystemInfoRequest,
&[],
MessageType::GetSystemInfoResponse,
)
}
pub async fn get_system_info(&mut self) -> Result<SystemInfo> {
self.unary_rpc(
MessageType::GetSystemInfoRequest,
&[],
MessageType::GetSystemInfoResponse,
)
.await
}
pub fn mmap_read_file_blocking(
&mut self,
path: &str,
offset: u64,
length: u64,
) -> Result<MmapReadFileResponse> {
let req = MmapReadFileRequest {
path: path.to_string(),
offset,
length,
..Default::default()
};
let payload = req.encode_to_vec();
self.unary_rpc_blocking(
MessageType::MmapReadFileRequest,
&payload,
MessageType::MmapReadFileResponse,
)
}
pub async fn ensure_runtime(&mut self, start_if_needed: bool) -> Result<RuntimeEnsureResponse> {
let req = RuntimeEnsureRequest {
start_if_needed,
..Default::default()
};
let payload = req.encode_to_vec();
self.unary_rpc(
MessageType::EnsureRuntimeRequest,
&payload,
MessageType::EnsureRuntimeResponse,
)
.await
}
pub fn ensure_runtime_blocking(
&mut self,
start_if_needed: bool,
) -> Result<RuntimeEnsureResponse> {
let req = RuntimeEnsureRequest {
start_if_needed,
..Default::default()
};
let payload = req.encode_to_vec();
let (resp_type, resp_payload) =
self.rpc_call_blocking(MessageType::EnsureRuntimeRequest, &payload)?;
if resp_type != MessageType::EnsureRuntimeResponse as u32 {
return Err(EngineError::Machine(format!(
"unexpected response type: {resp_type}"
)));
}
RuntimeEnsureResponse::decode_from_slice(&resp_payload)
.map_err(|e| EngineError::Machine(format!("failed to decode response: {e}")))
}
pub async fn get_runtime_status(&mut self) -> Result<RuntimeStatusResponse> {
let req = RuntimeStatusRequest::default();
let payload = req.encode_to_vec();
self.unary_rpc(
MessageType::RuntimeStatusRequest,
&payload,
MessageType::RuntimeStatusResponse,
)
.await
}
pub async fn watch_readiness(
mut self,
start_runtime_if_needed: bool,
timeout: Duration,
trace_id: &str,
) -> Result<ReadinessEvent> {
if !self.connected {
self.connect().await?;
}
self.require_agent_protocol().await?;
let req = WatchReadinessRequest {
start_runtime_if_needed,
timeout_ms: u32::try_from(timeout.as_millis()).unwrap_or(u32::MAX),
..Default::default()
};
let payload = req.encode_to_vec();
let buf = Self::build_message(MessageType::WatchReadinessRequest, trace_id, &payload);
match &mut self.transport {
AgentTransport::Async(t) => {
t.send(buf).await.map_err(|e| {
EngineError::Machine(format!("failed to send readiness watch request: {e}"))
})?;
let deadline = tokio::time::Instant::now() + timeout;
loop {
let remaining = deadline.saturating_duration_since(tokio::time::Instant::now());
if remaining.is_zero() {
return Err(EngineError::Machine(
"timeout waiting for guest readiness event".to_string(),
));
}
let raw = tokio::time::timeout(remaining, t.recv())
.await
.map_err(|_| {
EngineError::Machine(
"timeout waiting for guest readiness event".to_string(),
)
})?
.map_err(|e| {
EngineError::Machine(format!("failed to receive readiness event: {e}"))
})?;
let event = Self::decode_readiness_event(&raw)?;
if readiness_event_is_terminal(&event) || !start_runtime_if_needed {
return Ok(event);
}
}
}
AgentTransport::Blocking(t) => tokio::task::block_in_place(|| {
let deadline = Instant::now() + timeout;
t.send(&buf, deadline).map_err(|e| {
EngineError::Machine(format!("failed to send readiness watch request: {e}"))
})?;
loop {
let raw = t.recv(deadline).map_err(|e| {
EngineError::Machine(format!("failed to receive readiness event: {e}"))
})?;
let event = Self::decode_readiness_event(&raw)?;
if readiness_event_is_terminal(&event) || !start_runtime_if_needed {
return Ok(event);
}
}
}),
}
}
pub fn watch_readiness_blocking(
mut self,
start_runtime_if_needed: bool,
timeout: Duration,
trace_id: &str,
) -> Result<ReadinessEvent> {
self.require_agent_protocol_blocking()?;
let req = WatchReadinessRequest {
start_runtime_if_needed,
timeout_ms: u32::try_from(timeout.as_millis()).unwrap_or(u32::MAX),
..Default::default()
};
let payload = req.encode_to_vec();
let buf = Self::build_message(MessageType::WatchReadinessRequest, trace_id, &payload);
let AgentTransport::Blocking(t) = &mut self.transport else {
return Err(EngineError::Machine(
"blocking readiness watch called on async transport".to_string(),
));
};
let deadline = Instant::now() + timeout;
t.send(&buf, deadline).map_err(|e| {
EngineError::Machine(format!("failed to send readiness watch request: {e}"))
})?;
loop {
let raw = t.recv(deadline).map_err(|e| {
EngineError::Machine(format!("failed to receive readiness event: {e}"))
})?;
let event = Self::decode_readiness_event(&raw)?;
if readiness_event_is_terminal(&event) || !start_runtime_if_needed {
return Ok(event);
}
}
}
pub async fn watch_memory_pressure(&mut self, req: WatchMemoryPressureRequest) -> Result<()> {
if !self.connected {
self.connect().await?;
}
self.require_agent_protocol().await?;
let payload = req.encode_to_vec();
let buf = Self::build_message(MessageType::WatchMemoryPressureRequest, "", &payload);
match &mut self.transport {
AgentTransport::Async(t) => t.send(buf).await.map_err(|e| {
EngineError::Machine(format!("failed to send memory pressure watch request: {e}"))
}),
AgentTransport::Blocking(t) => tokio::task::block_in_place(|| {
let deadline = Instant::now() + Duration::from_secs(5);
t.send(&buf, deadline).map_err(|e| {
EngineError::Machine(format!(
"failed to send memory pressure watch request: {e}"
))
})
}),
}
}
pub async fn next_memory_pressure_event(
&mut self,
max_wait: Duration,
) -> Result<MemoryPressureEvent> {
match &mut self.transport {
AgentTransport::Async(t) => {
let raw = tokio::time::timeout(max_wait, t.recv())
.await
.map_err(|_| {
EngineError::Machine(
"timeout waiting for memory pressure event".to_string(),
)
})?
.map_err(|e| {
EngineError::Machine(format!(
"failed to receive memory pressure event: {e}"
))
})?;
Self::decode_memory_pressure_event(&raw)
}
AgentTransport::Blocking(t) => tokio::task::block_in_place(|| {
let raw = t.recv(Instant::now() + max_wait).map_err(|e| {
EngineError::Machine(format!("failed to receive memory pressure event: {e}"))
})?;
Self::decode_memory_pressure_event(&raw)
}),
}
}
pub async fn watch_stats(&mut self, req: WatchStatsRequest) -> Result<()> {
if !self.connected {
self.connect().await?;
}
self.require_agent_protocol().await?;
let payload = req.encode_to_vec();
let buf = Self::build_message(MessageType::WatchStatsRequest, "", &payload);
match &mut self.transport {
AgentTransport::Async(t) => t.send(buf).await.map_err(|e| {
EngineError::Machine(format!("failed to send stats watch request: {e}"))
}),
AgentTransport::Blocking(t) => tokio::task::block_in_place(|| {
let deadline = Instant::now() + Duration::from_secs(5);
t.send(&buf, deadline).map_err(|e| {
EngineError::Machine(format!("failed to send stats watch request: {e}"))
})
}),
}
}
pub async fn next_machine_stats(&mut self, max_wait: Duration) -> Result<MachineStats> {
match &mut self.transport {
AgentTransport::Async(t) => {
let raw = tokio::time::timeout(max_wait, t.recv())
.await
.map_err(|_| {
EngineError::Machine("timeout waiting for stats frame".to_string())
})?
.map_err(|e| {
EngineError::Machine(format!("failed to receive stats frame: {e}"))
})?;
Self::decode_machine_stats(&raw)
}
AgentTransport::Blocking(t) => tokio::task::block_in_place(|| {
let raw = t.recv(Instant::now() + max_wait).map_err(|e| {
EngineError::Machine(format!("failed to receive stats frame: {e}"))
})?;
Self::decode_machine_stats(&raw)
}),
}
}
fn decode_machine_stats(raw: &[u8]) -> Result<MachineStats> {
let (resp_type, _, resp_payload) = wire::parse_response(raw)?;
if resp_type == MessageType::Error as u32 {
let (code, message) = wire::parse_error_response(&resp_payload)?;
return Err(EngineError::Agent { code, message });
}
if resp_type != MessageType::MachineStats as u32 {
return Err(EngineError::Machine(format!(
"unexpected stats response type: 0x{resp_type:04x}"
)));
}
MachineStats::decode_from_slice(&resp_payload)
.map_err(|e| EngineError::Machine(format!("failed to decode machine stats: {e}")))
}
fn decode_memory_pressure_event(raw: &[u8]) -> Result<MemoryPressureEvent> {
let (resp_type, _, resp_payload) = wire::parse_response(raw)?;
if resp_type == MessageType::Error as u32 {
let (code, message) = wire::parse_error_response(&resp_payload)?;
return Err(EngineError::Agent { code, message });
}
if resp_type != MessageType::MemoryPressureEvent as u32 {
return Err(EngineError::Machine(format!(
"unexpected memory pressure response type: 0x{resp_type:04x}"
)));
}
MemoryPressureEvent::decode_from_slice(&resp_payload).map_err(|e| {
EngineError::Machine(format!("failed to decode memory pressure event: {e}"))
})
}
fn decode_readiness_event(raw: &[u8]) -> Result<ReadinessEvent> {
let (resp_type, _, resp_payload) = wire::parse_response(raw)?;
if resp_type == MessageType::Error as u32 {
let (code, message) = wire::parse_error_response(&resp_payload)?;
return Err(EngineError::Agent { code, message });
}
if resp_type != MessageType::ReadinessEvent as u32 {
return Err(EngineError::Machine(format!(
"unexpected readiness response type: 0x{resp_type:04x}"
)));
}
ReadinessEvent::decode_from_slice(&resp_payload)
.map_err(|e| EngineError::Machine(format!("failed to decode readiness event: {e}")))
}
pub async fn start_kubernetes(&mut self) -> Result<KubernetesStartResponse> {
let payload = KubernetesStartRequest::default().encode_to_vec();
self.unary_rpc(
MessageType::KubernetesStartRequest,
&payload,
MessageType::KubernetesStartResponse,
)
.await
}
pub async fn stop_kubernetes(&mut self) -> Result<KubernetesStopResponse> {
let payload = KubernetesStopRequest::default().encode_to_vec();
self.unary_rpc(
MessageType::KubernetesStopRequest,
&payload,
MessageType::KubernetesStopResponse,
)
.await
}
pub async fn delete_kubernetes(&mut self) -> Result<KubernetesDeleteResponse> {
let payload = KubernetesDeleteRequest::default().encode_to_vec();
self.unary_rpc(
MessageType::KubernetesDeleteRequest,
&payload,
MessageType::KubernetesDeleteResponse,
)
.await
}
pub async fn get_kubernetes_status(&mut self) -> Result<KubernetesStatusResponse> {
let payload = KubernetesStatusRequest::default().encode_to_vec();
self.unary_rpc(
MessageType::KubernetesStatusRequest,
&payload,
MessageType::KubernetesStatusResponse,
)
.await
}
pub async fn get_kubeconfig(&mut self) -> Result<KubernetesKubeconfigResponse> {
let payload = KubernetesKubeconfigRequest::default().encode_to_vec();
self.unary_rpc(
MessageType::KubernetesKubeconfigRequest,
&payload,
MessageType::KubernetesKubeconfigResponse,
)
.await
}
pub async fn list_kubernetes_load_balancers(
&mut self,
) -> Result<KubernetesLoadBalancersResponse> {
let payload = KubernetesLoadBalancersRequest::default().encode_to_vec();
self.unary_rpc(
MessageType::KubernetesLoadBalancersRequest,
&payload,
MessageType::KubernetesLoadBalancersResponse,
)
.await
}
pub async fn container_fs_paths(
&mut self,
container_id: &str,
) -> Result<ContainerFsPathsResponse> {
let payload = ContainerFsPathsRequest {
container_id: container_id.to_string(),
..Default::default()
}
.encode_to_vec();
self.unary_rpc(
MessageType::ContainerFsPathsRequest,
&payload,
MessageType::ContainerFsPathsResponse,
)
.await
}
pub fn container_fs_paths_blocking(
&mut self,
container_id: &str,
) -> Result<ContainerFsPathsResponse> {
let payload = ContainerFsPathsRequest {
container_id: container_id.to_string(),
..Default::default()
}
.encode_to_vec();
self.unary_rpc_blocking(
MessageType::ContainerFsPathsRequest,
&payload,
MessageType::ContainerFsPathsResponse,
)
}
pub async fn image_fs_paths(&mut self, top_chain_id: &str) -> Result<ImageFsPathsResponse> {
let payload = ImageFsPathsRequest {
top_chain_id: top_chain_id.to_string(),
..Default::default()
}
.encode_to_vec();
self.unary_rpc(
MessageType::ImageFsPathsRequest,
&payload,
MessageType::ImageFsPathsResponse,
)
.await
}
pub fn image_fs_paths_blocking(&mut self, top_chain_id: &str) -> Result<ImageFsPathsResponse> {
let payload = ImageFsPathsRequest {
top_chain_id: top_chain_id.to_string(),
..Default::default()
}
.encode_to_vec();
self.unary_rpc_blocking(
MessageType::ImageFsPathsRequest,
&payload,
MessageType::ImageFsPathsResponse,
)
}
pub async fn ensure_nfs_export(&mut self) -> Result<EnsureNfsExportResponse> {
let payload = EnsureNfsExportRequest::default().encode_to_vec();
self.unary_rpc(
MessageType::EnsureNfsExportRequest,
&payload,
MessageType::EnsureNfsExportResponse,
)
.await
}
pub fn ensure_nfs_export_blocking(&mut self) -> Result<EnsureNfsExportResponse> {
let payload = EnsureNfsExportRequest::default().encode_to_vec();
self.unary_rpc_blocking(
MessageType::EnsureNfsExportRequest,
&payload,
MessageType::EnsureNfsExportResponse,
)
}
pub async fn ensure_machine_export(
&mut self,
request: &EnsureMachineExportRequest,
) -> Result<EnsureMachineExportResponse> {
let payload = request.encode_to_vec();
self.unary_rpc(
MessageType::EnsureMachineExportRequest,
&payload,
MessageType::EnsureMachineExportResponse,
)
.await
}
pub fn ensure_machine_export_blocking(
&mut self,
request: &EnsureMachineExportRequest,
) -> Result<EnsureMachineExportResponse> {
let payload = request.encode_to_vec();
self.unary_rpc_blocking(
MessageType::EnsureMachineExportRequest,
&payload,
MessageType::EnsureMachineExportResponse,
)
}
pub async fn disk_trim(&mut self) -> Result<DiskTrimResponse> {
let payload = DiskTrimRequest::default().encode_to_vec();
self.unary_rpc(
MessageType::DiskTrimRequest,
&payload,
MessageType::DiskTrimResponse,
)
.await
}
pub fn disk_trim_blocking(&mut self) -> Result<DiskTrimResponse> {
let payload = DiskTrimRequest::default().encode_to_vec();
self.unary_rpc_blocking(
MessageType::DiskTrimRequest,
&payload,
MessageType::DiskTrimResponse,
)
}
pub async fn sandbox_create(
&mut self,
req: CreateSandboxRequest,
) -> Result<CreateSandboxResponse> {
let payload = req.encode_to_vec();
self.unary_rpc(
MessageType::SandboxCreateRequest,
&payload,
MessageType::SandboxCreateResponse,
)
.await
}
pub async fn sandbox_port_forward(
&mut self,
req: SandboxPortForwardRequest,
) -> Result<SandboxPortForwardResponse> {
let payload = req.encode_to_vec();
self.unary_rpc(
MessageType::SandboxPortForwardRequest,
&payload,
MessageType::SandboxPortForwardResponse,
)
.await
}
pub async fn sandbox_port_forward_remove(
&mut self,
req: SandboxPortForwardRemoveRequest,
) -> Result<()> {
let payload = req.encode_to_vec();
let (resp_type, _) = self
.rpc_call(MessageType::SandboxPortForwardRemoveRequest, &payload)
.await?;
Self::expect_ack_response_type(resp_type, MessageType::SandboxPortForwardRemoveResponse)
}
pub async fn sandbox_stop(
&mut self,
req: StopSandboxRequest,
) -> Result<SandboxCleanupResponse> {
let payload = req.encode_to_vec();
self.unary_rpc(
MessageType::SandboxStopRequest,
&payload,
MessageType::SandboxStopResponse,
)
.await
}
pub async fn sandbox_pause(
&mut self,
req: PauseSandboxRequest,
) -> Result<SandboxCleanupResponse> {
let payload = req.encode_to_vec();
self.unary_rpc(
MessageType::SandboxPauseRequest,
&payload,
MessageType::SandboxPauseResponse,
)
.await
}
pub async fn sandbox_resume(
&mut self,
req: SandboxResumeCommand,
) -> Result<SandboxResumeResponse> {
let payload = req.encode_to_vec();
self.unary_rpc(
MessageType::SandboxResumeRequest,
&payload,
MessageType::SandboxResumeResponse,
)
.await
}
pub async fn sandbox_set_lifecycle(&mut self, req: SetLifecycleRequest) -> Result<()> {
let (response_type, _) = self
.rpc_call(
MessageType::SandboxSetLifecycleRequest,
&req.encode_to_vec(),
)
.await?;
Self::expect_ack_response_type(response_type, MessageType::SandboxSetLifecycleResponse)
}
pub async fn sandbox_remove(
&mut self,
req: RemoveSandboxRequest,
) -> Result<SandboxCleanupResponse> {
let payload = req.encode_to_vec();
self.unary_rpc(
MessageType::SandboxRemoveRequest,
&payload,
MessageType::SandboxRemoveResponse,
)
.await
}
pub async fn sandbox_cleanup_prepare(&mut self, ticket: &SandboxCleanupTicket) -> Result<()> {
let (response_type, _) = self
.rpc_call(
MessageType::SandboxCleanupPrepareRequest,
&ticket.encode_to_vec(),
)
.await?;
Self::expect_ack_response_type(response_type, MessageType::SandboxCleanupPrepareResponse)
}
pub async fn sandbox_cleanup_finalize(&mut self, ticket: &SandboxCleanupTicket) -> Result<()> {
let (response_type, _) = self
.rpc_call(
MessageType::SandboxCleanupFinalizeRequest,
&ticket.encode_to_vec(),
)
.await?;
Self::expect_ack_response_type(response_type, MessageType::SandboxCleanupFinalizeResponse)
}
pub async fn sandbox_inspect(&mut self, req: InspectSandboxRequest) -> Result<SandboxInfo> {
let payload = req.encode_to_vec();
self.unary_rpc(
MessageType::SandboxInspectRequest,
&payload,
MessageType::SandboxInspectResponse,
)
.await
}
pub async fn sandbox_list(
&mut self,
req: ListSandboxesRequest,
) -> Result<ListSandboxesResponse> {
let payload = req.encode_to_vec();
self.unary_rpc(
MessageType::SandboxListRequest,
&payload,
MessageType::SandboxListResponse,
)
.await
}
pub async fn sandbox_read_file(self, req: ReadFileRequest) -> Result<SandboxStream<FileChunk>> {
self.open_sandbox_stream(
MessageType::SandboxFileReadRequest,
&req.encode_to_vec(),
StreamKind {
frame: MessageType::SandboxFileData,
end: None,
decode: |payload| {
FileChunk::decode_from_slice(payload).map_err(sandbox_stream::decode_error)
},
is_last: |chunk| chunk.done,
},
)
.await
}
pub async fn sandbox_write_file(
mut self,
open: WriteFileOpen,
mut data_rx: mpsc::Receiver<WriteFileChunk>,
) -> Result<()> {
if !self.connected {
self.connect().await?;
}
self.require_agent_protocol().await?;
let payload = open.encode_to_vec();
let buf = wire::build_message(MessageType::SandboxFileWriteRequest, "", &payload);
self.transport
.async_send(buf)
.await
.map_err(|e| EngineError::Machine(format!("failed to send write-file open: {e}")))?;
while let Some(item) = data_rx.recv().await {
let data = match item {
WriteFileChunk::Data(data) => data,
WriteFileChunk::Abort => {
return Err(EngineError::Machine(
"write_file aborted: client stream ended before completion".to_owned(),
));
}
};
let chunk = FileChunk {
data,
done: false,
..Default::default()
};
let frame =
wire::build_message(MessageType::SandboxFileChunk, "", &chunk.encode_to_vec());
self.transport
.async_send(frame)
.await
.map_err(|e| EngineError::Machine(format!("failed to send file chunk: {e}")))?;
}
let done = FileChunk {
data: Vec::new(),
done: true,
..Default::default()
};
let frame = wire::build_message(MessageType::SandboxFileChunk, "", &done.encode_to_vec());
self.transport
.async_send(frame)
.await
.map_err(|e| EngineError::Machine(format!("failed to send final chunk: {e}")))?;
let raw = self
.transport
.async_recv()
.await
.map_err(|e| EngineError::Machine(format!("recv error: {e}")))?;
let (resp_type, _, resp_payload) = wire::parse_response(&raw)?;
if resp_type == MessageType::Error as u32 {
let (code, message) = wire::parse_error_response(&resp_payload)?;
return Err(EngineError::Agent { code, message });
}
if resp_type != MessageType::SandboxFileWriteResponse as u32 {
return Err(EngineError::Machine(format!(
"unexpected response type: 0x{resp_type:04x}"
)));
}
Ok(())
}
pub async fn sandbox_stat(&mut self, req: StatFileRequest) -> Result<FileStat> {
let payload = req.encode_to_vec();
self.unary_rpc(
MessageType::SandboxFileStatRequest,
&payload,
MessageType::SandboxFileStatResponse,
)
.await
}
pub async fn sandbox_list_dir(&mut self, req: ListDirRequest) -> Result<ListDirResponse> {
let payload = req.encode_to_vec();
self.unary_rpc(
MessageType::SandboxFileListDirRequest,
&payload,
MessageType::SandboxFileListDirResponse,
)
.await
}
pub async fn sandbox_make_dir(&mut self, req: MakeDirRequest) -> Result<()> {
let (resp_type, _) = self
.rpc_call(MessageType::SandboxFileMakeDirRequest, &req.encode_to_vec())
.await?;
Self::expect_ack_response_type(resp_type, MessageType::SandboxFileMakeDirResponse)
}
pub async fn sandbox_remove_entry(&mut self, req: RemoveEntryRequest) -> Result<()> {
let (resp_type, _) = self
.rpc_call(MessageType::SandboxFileRemoveRequest, &req.encode_to_vec())
.await?;
Self::expect_ack_response_type(resp_type, MessageType::SandboxFileRemoveResponse)
}
pub async fn sandbox_move_entry(&mut self, req: MoveEntryRequest) -> Result<()> {
let (resp_type, _) = self
.rpc_call(MessageType::SandboxFileMoveRequest, &req.encode_to_vec())
.await?;
Self::expect_ack_response_type(resp_type, MessageType::SandboxFileMoveResponse)
}
pub async fn sandbox_watch_dir(
self,
req: WatchDirRequest,
) -> Result<SandboxStream<WatchDirResponse>> {
self.open_sandbox_stream(
MessageType::SandboxFileWatchRequest,
&req.encode_to_vec(),
StreamKind {
frame: MessageType::SandboxFileWatchEvent,
end: Some(MessageType::SandboxFileWatchEnd),
decode: |payload| {
WatchDirResponse::decode_from_slice(payload)
.map_err(sandbox_stream::decode_error)
},
is_last: |_| false,
},
)
.await
}
pub async fn sandbox_exec_start(&mut self, req: StartExecutionRequest) -> Result<Execution> {
let payload = req.encode_to_vec();
self.unary_rpc(
MessageType::SandboxExecStartRequest,
&payload,
MessageType::SandboxExecStartResponse,
)
.await
}
pub async fn sandbox_stdin_write(&mut self, req: WriteStdinRequest) -> Result<StdinStatus> {
let payload = req.encode_to_vec();
self.unary_rpc(
MessageType::SandboxStdinWriteRequest,
&payload,
MessageType::SandboxStdinStatus,
)
.await
}
pub async fn sandbox_stdin_status(
&mut self,
req: GetStdinStatusRequest,
) -> Result<StdinStatus> {
let payload = req.encode_to_vec();
self.unary_rpc(
MessageType::SandboxStdinStatusRequest,
&payload,
MessageType::SandboxStdinStatus,
)
.await
}
pub async fn sandbox_exec_signal(&mut self, req: SignalExecutionRequest) -> Result<()> {
let payload = req.encode_to_vec();
let (resp_type, _) = self
.rpc_call(MessageType::SandboxExecSignalRequest, &payload)
.await?;
Self::expect_ack_response_type(resp_type, MessageType::SandboxExecSignalResponse)
}
pub async fn sandbox_exec_resize(&mut self, req: ResizeExecutionTtyRequest) -> Result<()> {
let payload = req.encode_to_vec();
let (resp_type, _) = self
.rpc_call(MessageType::SandboxExecResizeRequest, &payload)
.await?;
Self::expect_ack_response_type(resp_type, MessageType::SandboxExecResizeResponse)
}
pub async fn sandbox_exec_wait(&mut self, req: WaitExecutionRequest) -> Result<Execution> {
let payload = req.encode_to_vec();
self.unary_rpc(
MessageType::SandboxExecWaitRequest,
&payload,
MessageType::SandboxExecWaitResponse,
)
.await
}
pub async fn sandbox_exec_list(
&mut self,
req: ListExecutionsRequest,
) -> Result<ListExecutionsResponse> {
let payload = req.encode_to_vec();
self.unary_rpc(
MessageType::SandboxExecListRequest,
&payload,
MessageType::SandboxExecListResponse,
)
.await
}
pub async fn sandbox_wait_for_port(&mut self, req: WaitForPortRequest) -> Result<()> {
let (resp_type, _) = self
.rpc_call(MessageType::SandboxWaitForPortRequest, &req.encode_to_vec())
.await?;
Self::expect_ack_response_type(resp_type, MessageType::SandboxWaitForPortResponse)
}
pub async fn sandbox_exec_attach(
self,
req: AttachExecutionRequest,
) -> Result<SandboxStream<ExecutionEvent>> {
self.open_sandbox_stream(
MessageType::SandboxExecAttachRequest,
&req.encode_to_vec(),
StreamKind {
frame: MessageType::SandboxExecEvent,
end: None,
decode: |payload| {
ExecutionEvent::decode_from_slice(payload).map_err(sandbox_stream::decode_error)
},
is_last: |event| matches!(event.event, Some(execution_event::Event::Exited(_))),
},
)
.await
}
pub async fn sandbox_events(
self,
req: SandboxEventsRequest,
) -> Result<SandboxStream<SandboxEvent>> {
self.open_sandbox_stream(
MessageType::SandboxEventsRequest,
&req.encode_to_vec(),
StreamKind {
frame: MessageType::SandboxEvent,
end: None,
decode: |payload| {
SandboxEvent::decode_from_slice(payload).map_err(sandbox_stream::decode_error)
},
is_last: |_| false,
},
)
.await
}
pub async fn sandbox_cleanup_events(self) -> Result<SandboxStream<SandboxCleanupTicket>> {
self.open_sandbox_stream(
MessageType::WatchSandboxCleanupRequest,
&WatchSandboxCleanupRequest::default().encode_to_vec(),
StreamKind {
frame: MessageType::SandboxCleanupEvent,
end: None,
decode: |payload| {
SandboxCleanupTicket::decode_from_slice(payload)
.map_err(sandbox_stream::decode_error)
},
is_last: |_| false,
},
)
.await
}
pub async fn sandbox_checkpoint(
&mut self,
req: CheckpointRequest,
) -> Result<CheckpointResponse> {
let payload = req.encode_to_vec();
self.unary_rpc(
MessageType::SandboxCheckpointRequest,
&payload,
MessageType::SandboxCheckpointResponse,
)
.await
}
pub async fn sandbox_restore(&mut self, req: RestoreRequest) -> Result<RestoreResponse> {
let payload = req.encode_to_vec();
self.unary_rpc(
MessageType::SandboxRestoreRequest,
&payload,
MessageType::SandboxRestoreResponse,
)
.await
}
pub async fn sandbox_list_snapshots(
&mut self,
req: ListSnapshotsRequest,
) -> Result<ListSnapshotsResponse> {
let payload = req.encode_to_vec();
self.unary_rpc(
MessageType::SandboxListSnapshotsRequest,
&payload,
MessageType::SandboxListSnapshotsResponse,
)
.await
}
pub async fn sandbox_delete_snapshot(&mut self, req: DeleteSnapshotRequest) -> Result<()> {
let payload = req.encode_to_vec();
let (resp_type, _) = self
.rpc_call(MessageType::SandboxDeleteSnapshotRequest, &payload)
.await?;
Self::expect_ack_response_type(resp_type, MessageType::SandboxDeleteSnapshotResponse)
}
pub async fn sandbox_template_build(&mut self, req: BuildTemplateRequest) -> Result<Template> {
let payload = req.encode_to_vec();
self.unary_rpc(
MessageType::SandboxTemplateBuildRequest,
&payload,
MessageType::SandboxTemplateBuildResponse,
)
.await
}
pub async fn sandbox_template_publish(
&mut self,
req: PublishTemplateRequest,
) -> Result<Template> {
let payload = req.encode_to_vec();
self.unary_rpc(
MessageType::SandboxTemplatePublishRequest,
&payload,
MessageType::SandboxTemplatePublishResponse,
)
.await
}
pub async fn sandbox_template_get(&mut self, req: GetTemplateRequest) -> Result<Template> {
let payload = req.encode_to_vec();
self.unary_rpc(
MessageType::SandboxTemplateGetRequest,
&payload,
MessageType::SandboxTemplateGetResponse,
)
.await
}
pub async fn sandbox_template_list(
&mut self,
req: ListTemplatesRequest,
) -> Result<ListTemplatesResponse> {
let payload = req.encode_to_vec();
self.unary_rpc(
MessageType::SandboxTemplateListRequest,
&payload,
MessageType::SandboxTemplateListResponse,
)
.await
}
pub async fn sandbox_template_delete(&mut self, req: DeleteTemplateRequest) -> Result<()> {
let payload = req.encode_to_vec();
let (resp_type, _) = self
.rpc_call(MessageType::SandboxTemplateDeleteRequest, &payload)
.await?;
Self::expect_ack_response_type(resp_type, MessageType::SandboxTemplateDeleteResponse)
}
}
fn readiness_event_is_terminal(event: &ReadinessEvent) -> bool {
use arcbox_connect::v1::readiness_event::Kind;
matches!(
event.kind.as_known(),
Some(Kind::RuntimeReady | Kind::RuntimeFailed)
)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_message_type_roundtrip() {
assert_eq!(
MessageType::from_u32(MessageType::PingRequest as u32),
Some(MessageType::PingRequest)
);
assert_eq!(
MessageType::from_u32(MessageType::PingResponse as u32),
Some(MessageType::PingResponse)
);
assert_eq!(
MessageType::from_u32(MessageType::PortBindingsChanged as u32),
Some(MessageType::PortBindingsChanged)
);
}
#[test]
fn test_agent_client_new() {
let client = AgentClient::new(3);
assert_eq!(client.cid(), 3);
assert!(!client.connected);
}
fn ping_response(protocol_version: u32) -> PingResponse {
PingResponse {
message: "pong".to_string(),
version: "0.4.16".to_string(),
protocol_version,
..Default::default()
}
}
#[test]
fn boot_contract_rejection_is_not_reported_as_a_stale_agent() {
let mut response = ping_response(0);
response.message =
"incompatible host boot contract: missing arcbox.runtime_generation".to_string();
let err =
AgentClient::check_agent_protocol(&response).expect_err("protocol 0 must be rejected");
assert!(err.to_string().contains("incompatible"));
assert!(err.to_string().contains("arcbox.runtime_generation"));
assert!(!err.to_string().contains("stale"));
}
#[test]
fn previous_protocol_is_rejected() {
let previous = arcbox_constants::wire::AGENT_PROTOCOL_VERSION - 1;
let err = AgentClient::check_agent_protocol(&ping_response(previous))
.expect_err("previous protocol must be rejected");
assert!(err.to_string().contains("incompatible"));
assert!(err.to_string().contains("stale"));
}
#[test]
fn current_protocol_is_accepted() {
let resp = ping_response(arcbox_constants::wire::AGENT_PROTOCOL_VERSION);
AgentClient::check_agent_protocol(&resp).expect("current protocol must pass");
}
#[test]
fn newer_agent_protocol_is_accepted_with_warning() {
let resp = ping_response(arcbox_constants::wire::AGENT_PROTOCOL_VERSION + 1);
AgentClient::check_agent_protocol(&resp).expect("newer protocol must pass");
}
}