use std::collections::{BTreeMap, VecDeque};
use std::sync::Arc;
use std::time::Duration;
use async_trait::async_trait;
use base64::engine::general_purpose::STANDARD as BASE64;
use base64::Engine as _;
use futures::stream::{self, BoxStream};
use serde::Deserialize;
use serde_json::json;
use tracing::warn;
use crate::error::{ErrorData, Result};
use crate::traits::{
Binding, CommandOutput, CreateSandboxRequest, JobError, JobExit, JobPoll, JobStart,
PreviewCapability, ResolvedSandbox, RunCommandRequest, Sandbox, SandboxInstance, SandboxState,
};
use alien_core::{SandboxCapabilities, SandboxEgress};
use alien_error::{AlienError, Context, ContextError};
use alien_gcp_clients::gcp::agent_platform::{
AgentPlatformApi, AgentPlatformErrorData, EgressControlConfig, SandboxCreateRequest,
SandboxEnvironment, SandboxSnapshot,
};
use alien_gcp_clients::gcp::longrunning::{Operation, OperationResult};
const AGENT_PROTOCOL_VERSION: u32 = 1;
const MAX_SYNCHRONOUS_TIMEOUT: Duration = Duration::from_secs(30);
const MAX_SANDBOX_ID: usize = 63;
const SANDBOX_READY_ATTEMPTS: u32 = 150;
const SANDBOX_READY_INTERVAL: Duration = Duration::from_secs(2);
#[cfg(not(test))]
const AGENT_READY_TIMEOUT: Duration = Duration::from_secs(60);
#[cfg(not(test))]
const AGENT_READY_POLL: Duration = Duration::from_secs(1);
#[cfg(test)]
const AGENT_READY_TIMEOUT: Duration = Duration::from_millis(50);
#[cfg(test)]
const AGENT_READY_POLL: Duration = Duration::from_millis(5);
const OPERATION_POLL_ATTEMPTS: u32 = 150;
const OPERATION_POLL_INTERVAL: Duration = Duration::from_secs(2);
const TERMINATE_POLL_ATTEMPTS: u32 = 30;
const TERMINATE_POLL_INTERVAL: Duration = Duration::from_secs(2);
const JOB_POLL_INTERVAL: Duration = Duration::from_secs(1);
const JOB_POLL_GRACE: Duration = Duration::from_secs(15);
const CREATE: &str = "sandbox.create";
const GET: &str = "sandbox.get";
const GET_OR_CREATE: &str = "sandbox.getOrCreate";
const RUN_COMMAND: &str = "sandbox.runCommand";
const JOB_START: &str = "sandbox.jobStart";
const JOB_POLL: &str = "sandbox.jobPoll";
const JOB_CANCEL: &str = "sandbox.jobCancel";
const TERMINATE: &str = "sandbox.terminate";
const NO_GENERATION: u64 = 0;
const AGENT_PROBE_BUDGET: Duration = Duration::from_secs(60);
pub fn egress_control_config(
sandbox_label: &str,
egress: &SandboxEgress,
) -> Result<EgressControlConfig> {
let Some(internet_access) = egress.internet_access_switch() else {
return Err(AlienError::new(ErrorData::InvalidInput {
operation_context: "sandbox.template".to_string(),
details: format!(
"sandbox '{sandbox_label}' asked for domain-scoped egress, which Agent \
Platform cannot express; it offers only 'allow' (open) and 'deny' (closed)"
),
field_name: Some("egress".to_string()),
}));
};
Ok(EgressControlConfig {
internet_access: Some(internet_access),
extra: Default::default(),
})
}
#[derive(Debug)]
pub struct GcpAgentPlatformSandbox {
client: Arc<dyn AgentPlatformApi>,
engine: String,
template: String,
max_lifetime_seconds: Option<u32>,
}
impl GcpAgentPlatformSandbox {
fn lifetime_seconds(&self, timeout_ms: Option<u64>, operation: &str) -> Result<Option<u32>> {
match timeout_ms {
Some(timeout_ms) => {
super::requested_lifetime_seconds(timeout_ms, self.max_lifetime_seconds, operation)
.map(Some)
}
None => Ok(self.max_lifetime_seconds),
}
}
pub fn new(
client: Arc<dyn AgentPlatformApi>,
engine: String,
template: String,
max_lifetime_seconds: Option<u32>,
) -> Self {
let engine = engine.rsplit('/').next().unwrap_or(&engine).to_string();
Self {
client,
engine,
template,
max_lifetime_seconds,
}
}
#[cfg(test)]
pub(crate) fn engine(&self) -> &str {
&self.engine
}
fn unsupported(&self, capability: &str, reason: &str) -> AlienError<ErrorData> {
AlienError::new(ErrorData::OperationNotSupported {
operation: capability.to_string(),
reason: reason.to_string(),
})
}
fn pause_resume_unsupported(&self) -> AlienError<ErrorData> {
self.unsupported(
alien_core::SandboxCapability::PauseResume.as_str(),
"Agent Platform sandboxes cannot be paused and resumed with their state kept",
)
}
fn checked_sandbox_id(operation: &str, sandbox_id: &str) -> Result<()> {
if is_addressable_id(sandbox_id) {
return Ok(());
}
Err(AlienError::new(ErrorData::InvalidInput {
operation_context: operation.to_string(),
details: format!(
"sandbox id '{sandbox_id}' must be a single segment of letters, digits, '-' and \
'_', at most {MAX_SANDBOX_ID} characters"
),
field_name: Some("sandboxId".to_string()),
}))
}
async fn read_sandbox(
&self,
operation: &str,
sandbox_id: &str,
) -> Result<Option<SandboxEnvironment>> {
match self.client.get_sandbox(&self.engine, sandbox_id).await {
Ok(sandbox) => Ok(Some(sandbox)),
Err(error) if is_not_found(&error) => Ok(None),
Err(error) => Err(error.context(ErrorData::SandboxUnreachable {
operation: operation.to_string(),
reason: "the Agent Platform API did not answer a sandbox read".to_string(),
})),
}
}
async fn await_operation(
&self,
operation: &str,
started: Operation,
) -> Result<serde_json::Value> {
let Some(name) = started.name.clone() else {
return Err(AlienError::new(ErrorData::UnexpectedResponseFormat {
provider: "gcp-agent-platform".to_string(),
binding_name: operation.to_string(),
field: "name".to_string(),
response_json: "the operation carried no resource name to poll".to_string(),
}));
};
let mut current = started;
for _ in 0..OPERATION_POLL_ATTEMPTS {
if current.done == Some(true) {
return finish_operation(operation, &name, current);
}
tokio::time::sleep(OPERATION_POLL_INTERVAL).await;
current =
self.client
.get_operation(&name)
.await
.context(ErrorData::SandboxUnreachable {
operation: operation.to_string(),
reason: format!("could not read operation '{name}'"),
})?;
}
if current.done == Some(true) {
return finish_operation(operation, &name, current);
}
Err(AlienError::new(ErrorData::SandboxUnreachable {
operation: operation.to_string(),
reason: format!("operation '{name}' did not complete within its polling budget"),
}))
}
async fn execute_op(
&self,
sandbox_id: &str,
operation: &str,
envelope: serde_json::Value,
) -> Result<Vec<u8>> {
let body = serde_json::to_vec(&envelope).map_err(|error| {
AlienError::new(ErrorData::SerializationFailed {
message: format!("could not encode the {operation} envelope: {error}"),
})
})?;
self.client
.execute(&self.engine, sandbox_id, &body)
.await
.map_err(|error| Self::execute_failed(operation, error))
}
fn execute_failed(
operation: &str,
error: AlienError<AgentPlatformErrorData>,
) -> AlienError<ErrorData> {
if let Some(answer) = agent_answer(&error) {
return error.context(ErrorData::SandboxCommandFailed {
failure: "agentRefused".to_string(),
reason: format!("{operation} was refused: {answer}"),
});
}
if is_not_found(&error) {
return error.context(ErrorData::SandboxCommandFailed {
failure: "sandboxGone".to_string(),
reason: format!("{operation}: the sandbox does not exist"),
});
}
if operation == RUN_COMMAND || operation == JOB_START {
return error.context(ErrorData::SandboxOutcomeUnknown {
operation: operation.to_string(),
reason: "the sandbox did not complete the call".to_string(),
});
}
error.context(ErrorData::SandboxCommandFailed {
failure: "executeFailed".to_string(),
reason: format!("{operation} could not be completed against the sandbox"),
})
}
async fn probe_agent(&self, operation: &str, sandbox_id: &str) -> Result<u64> {
let unreachable = |reason: String| {
AlienError::new(ErrorData::SandboxUnreachable {
operation: operation.to_string(),
reason,
})
};
let body = tokio::time::timeout(
AGENT_PROBE_BUDGET,
self.client.execute(
&self.engine,
sandbox_id,
&serde_json::to_vec(&json!({ "v": AGENT_PROTOCOL_VERSION, "op": "health" }))
.unwrap_or_default(),
),
)
.await
.map_err(|_| {
unreachable(format!(
"the sandbox's agent did not answer a health probe within {}s",
AGENT_PROBE_BUDGET.as_secs()
))
})?
.map_err(|error| {
error.context(ErrorData::SandboxUnreachable {
operation: operation.to_string(),
reason: "the sandbox's agent did not answer a health probe".to_string(),
})
})?;
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
struct Health {
protocol_version: u32,
boot_id: String,
}
let health: Health = serde_json::from_slice(&body).map_err(|_| {
unreachable(format!(
"the sandbox's agent answered a health probe with a body this provider cannot \
read: {}",
truncated(&body)
))
})?;
if health.protocol_version != AGENT_PROTOCOL_VERSION {
return Err(unreachable(format!(
"the sandbox's agent speaks protocol {} where this provider speaks {}",
health.protocol_version, AGENT_PROTOCOL_VERSION
)));
}
if health.boot_id.is_empty() {
return Err(unreachable(
"the sandbox's agent reported no container boot id, so its identity cannot be \
established"
.to_string(),
));
}
Ok(generation_from_boot_id(&health.boot_id))
}
async fn discard(
&self,
sandbox_id: &str,
reason: AlienError<ErrorData>,
) -> AlienError<ErrorData> {
let Err(error) = self.client.delete_sandbox(&self.engine, sandbox_id).await else {
return reason;
};
warn!(
sandbox = %sandbox_id,
%error,
"could not delete a sandbox that was never handed to its caller"
);
reason.context(ErrorData::SandboxCommandFailed {
failure: "sandboxLeftBehind".to_string(),
reason: format!(
"sandbox '{sandbox_id}' was not handed to its caller and could not be deleted, so \
it is still running"
),
})
}
async fn settle(&self, sandbox_id: &str) -> Result<u64> {
for _ in 0..SANDBOX_READY_ATTEMPTS {
let Some(sandbox) = self.read_sandbox(CREATE, sandbox_id).await? else {
return Err(AlienError::new(ErrorData::SandboxCommandFailed {
failure: "sandboxGone".to_string(),
reason: format!("sandbox '{sandbox_id}' disappeared while it was coming up"),
}));
};
match sandbox_state(CREATE, sandbox.state.as_deref())? {
SandboxState::Running => return self.wait_until_servable(sandbox_id).await,
SandboxState::Terminated => {
return Err(AlienError::new(ErrorData::SandboxCommandFailed {
failure: "sandboxTerminated".to_string(),
reason: format!(
"sandbox '{sandbox_id}' reached a terminal state while starting"
),
}));
}
SandboxState::Starting | SandboxState::Paused => {}
}
tokio::time::sleep(SANDBOX_READY_INTERVAL).await;
}
Err(AlienError::new(ErrorData::SandboxUnreachable {
operation: CREATE.to_string(),
reason: format!(
"sandbox '{sandbox_id}' was not running after {}s",
SANDBOX_READY_ATTEMPTS as u64 * SANDBOX_READY_INTERVAL.as_secs()
),
}))
}
async fn wait_until_servable(&self, sandbox_id: &str) -> Result<u64> {
let deadline = tokio::time::Instant::now() + AGENT_READY_TIMEOUT;
let mut last_error = None;
loop {
match tokio::time::timeout_at(deadline, self.probe_agent(CREATE, sandbox_id)).await {
Ok(Ok(generation)) => return Ok(generation),
Ok(Err(error)) => {
last_error = Some(error);
if tokio::time::Instant::now() + AGENT_READY_POLL >= deadline {
break;
}
tokio::time::sleep(AGENT_READY_POLL).await;
}
Err(_) => break,
}
}
let timed_out = ErrorData::SandboxUnreachable {
operation: CREATE.to_string(),
reason: format!(
"sandbox '{sandbox_id}' was running but its agent did not become servable within {}s",
AGENT_READY_TIMEOUT.as_secs()
),
};
Err(match last_error {
Some(error) => error.context(timed_out),
None => AlienError::new(timed_out),
})
}
async fn run_synchronous(
&self,
sandbox_id: &str,
request: &RunCommandRequest,
) -> Result<BoxStream<'static, Result<CommandOutput>>> {
let envelope = exec_envelope("exec", sandbox_id, request);
let body = self.execute_op(sandbox_id, RUN_COMMAND, envelope).await?;
let frames = parse_exec_frames(&body)?;
Ok(Box::pin(stream::iter(frames)))
}
async fn run_detached(
&self,
sandbox_id: &str,
request: RunCommandRequest,
) -> Result<BoxStream<'static, Result<CommandOutput>>> {
let timeout = request.timeout;
let started = self.start_job(sandbox_id, request).await?;
let state = JobPollState {
client: self.client.clone(),
engine: self.engine.clone(),
sandbox_id: sandbox_id.to_string(),
job_id: started.job_id,
since_seq: None,
pending: VecDeque::new(),
finished: false,
stopped: false,
deadline_at: tokio::time::Instant::now() + timeout + JOB_POLL_GRACE,
};
Ok(Box::pin(stream::unfold(state, job_poll_step)))
}
fn checked_command(operation: &str, request: &RunCommandRequest) -> Result<()> {
if request.command.is_empty() {
return Err(AlienError::new(ErrorData::InvalidInput {
operation_context: operation.to_string(),
details: "a command must name a program to run".to_string(),
field_name: Some("command".to_string()),
}));
}
if timeout_millis(request.timeout) == 0 {
return Err(AlienError::new(ErrorData::SandboxCommandFailed {
failure: "invalidRequest".to_string(),
reason: "a command must carry a timeout of at least one millisecond".to_string(),
}));
}
Ok(())
}
}
impl Binding for GcpAgentPlatformSandbox {}
#[async_trait]
impl Sandbox for GcpAgentPlatformSandbox {
fn as_any(&self) -> &dyn std::any::Any {
self
}
fn capabilities(&self) -> SandboxCapabilities {
SandboxCapabilities::gcp_agent_platform()
}
async fn create(&self, request: CreateSandboxRequest) -> Result<SandboxInstance> {
if !request.env.is_empty() {
return Err(AlienError::new(ErrorData::OperationNotSupported {
operation: CREATE.to_string(),
reason: "Agent Platform sandboxes take no sandbox-level env; set env per command \
instead"
.to_string(),
}));
}
if request.tenant_key.is_some() {
return Err(AlienError::new(ErrorData::OperationNotSupported {
operation: CREATE.to_string(),
reason: "Agent Platform sandboxes take no tenantKey; create one sandbox per \
tenant instead"
.to_string(),
}));
}
let ttl = self
.lifetime_seconds(request.timeout_ms, CREATE)?
.map(|seconds| format!("{seconds}s"));
let started = self
.client
.create_sandbox(
&self.engine,
SandboxCreateRequest {
display_name: request.sandbox_id.clone(),
sandbox_environment_template: Some(self.template.clone()),
sandbox_environment_snapshot: None,
ttl,
},
)
.await
.context(ErrorData::SandboxUnreachable {
operation: CREATE.to_string(),
reason: "the Agent Platform API refused a sandbox create".to_string(),
})?;
let created: SandboxEnvironment = serde_json::from_value(
self.await_operation(CREATE, started).await?,
)
.map_err(|error| {
AlienError::new(ErrorData::UnexpectedResponseFormat {
provider: "gcp-agent-platform".to_string(),
binding_name: CREATE.to_string(),
field: "response".to_string(),
response_json: format!("the create operation resolved to a non-sandbox: {error}"),
})
})?;
let Some(sandbox_id) = created.name.as_deref().and_then(sandbox_segment) else {
return Err(AlienError::new(ErrorData::UnexpectedResponseFormat {
provider: "gcp-agent-platform".to_string(),
binding_name: CREATE.to_string(),
field: "name".to_string(),
response_json: format!("{:?}", created.name),
}));
};
let sandbox_id = sandbox_id.to_string();
match self.settle(&sandbox_id).await {
Ok(generation) => Ok(SandboxInstance {
sandbox_id,
state: SandboxState::Running,
generation,
}),
Err(error) => Err(self.discard(&sandbox_id, error).await),
}
}
async fn get(&self, sandbox_id: &str) -> Result<Option<SandboxInstance>> {
Self::checked_sandbox_id(GET, sandbox_id)?;
let Some(sandbox) = self.read_sandbox(GET, sandbox_id).await? else {
return Ok(None);
};
let state = sandbox_state(GET, sandbox.state.as_deref())?;
let generation = if state == SandboxState::Running {
self.probe_agent(GET, sandbox_id).await?
} else {
NO_GENERATION
};
Ok(Some(SandboxInstance {
sandbox_id: sandbox_id.to_string(),
state,
generation,
}))
}
async fn get_or_create(&self, request: CreateSandboxRequest) -> Result<ResolvedSandbox> {
if let Some(id) = request.sandbox_id.as_deref() {
match self.get(id).await {
Ok(Some(sandbox)) if sandbox.state == SandboxState::Running => {
return Ok(ResolvedSandbox::found(sandbox))
}
Ok(Some(sandbox)) if sandbox.state == SandboxState::Starting => {
let generation = self.settle(id).await?;
return Ok(ResolvedSandbox::found(SandboxInstance {
sandbox_id: id.to_string(),
state: SandboxState::Running,
generation,
}));
}
Ok(_) => {}
Err(error) if error.code == "SANDBOX_UNREACHABLE" => {}
Err(error) => {
return Err(error.context(ErrorData::SandboxCommandFailed {
failure: "getOrCreateFailed".to_string(),
reason: format!("{GET_OR_CREATE}: reaching sandbox '{id}' failed"),
}))
}
}
}
self.create(request).await.map(ResolvedSandbox::created)
}
async fn list(&self) -> Result<Vec<SandboxInstance>> {
let sandboxes = self.client.list_sandboxes(&self.engine).await.context(
ErrorData::SandboxUnreachable {
operation: "sandbox.list".to_string(),
reason: "the Agent Platform API did not answer a sandbox list".to_string(),
},
)?;
Ok(sandboxes
.into_iter()
.filter_map(|sandbox| {
let sandbox_id = sandbox.name.as_deref().and_then(sandbox_segment)?;
let state = sandbox_state("sandbox.list", sandbox.state.as_deref()).ok()?;
Some(SandboxInstance {
sandbox_id: sandbox_id.to_string(),
state,
generation: NO_GENERATION,
})
})
.collect())
}
async fn run_command(
&self,
sandbox_id: &str,
request: RunCommandRequest,
) -> Result<BoxStream<'static, Result<CommandOutput>>> {
Self::checked_sandbox_id(RUN_COMMAND, sandbox_id)?;
Self::checked_command(RUN_COMMAND, &request)?;
if request.timeout <= MAX_SYNCHRONOUS_TIMEOUT {
self.run_synchronous(sandbox_id, &request).await
} else {
self.run_detached(sandbox_id, request).await
}
}
async fn start_job(&self, sandbox_id: &str, request: RunCommandRequest) -> Result<JobStart> {
Self::checked_sandbox_id(JOB_START, sandbox_id)?;
Self::checked_command(JOB_START, &request)?;
let envelope = exec_envelope("jobStart", sandbox_id, &request);
let body = self.execute_op(sandbox_id, JOB_START, envelope).await?;
let started: JobStartReply = serde_json::from_slice(&body).map_err(|_| {
AlienError::new(ErrorData::UnexpectedResponseFormat {
provider: "gcp-agent-platform".to_string(),
binding_name: JOB_START.to_string(),
field: "jobId".to_string(),
response_json: truncated(&body),
})
.context(ErrorData::SandboxOutcomeUnknown {
operation: JOB_START.to_string(),
reason: "the job started and its id could not be read, so it cannot be polled"
.to_string(),
})
})?;
Ok(JobStart {
job_id: started.job_id,
})
}
async fn poll_job(
&self,
sandbox_id: &str,
job_id: &str,
since_seq: Option<u64>,
) -> Result<JobPoll> {
Self::checked_sandbox_id(JOB_POLL, sandbox_id)?;
let reply = poll_once(
self.client.as_ref(),
&self.engine,
sandbox_id,
job_id,
since_seq,
)
.await?;
Ok(JobPoll {
running: reply.running,
frames: reply
.frames
.into_iter()
.map(WireFrame::into_output)
.collect::<Result<Vec<_>>>()?,
exit: reply.exit_code.map(|code| JobExit {
code,
truncated: reply.truncated.unwrap_or(false),
}),
error: reply.error.map(|error| JobError {
code: error.code,
message: error.message,
}),
})
}
async fn cancel_job(&self, sandbox_id: &str, job_id: &str) -> Result<()> {
Self::checked_sandbox_id(JOB_CANCEL, sandbox_id)?;
let body = self
.client
.execute(&self.engine, sandbox_id, &cancel_body(job_id))
.await
.map_err(|error| unanswered_job(JOB_CANCEL, error))?;
if !cancel_confirmed(&body) {
return Err(AlienError::new(ErrorData::SandboxCommandFailed {
failure: "agentRefused".to_string(),
reason: format!("{JOB_CANCEL}: {}", truncated(&body)),
}));
}
Ok(())
}
async fn read_file(&self, sandbox_id: &str, path: &str) -> Result<Vec<u8>> {
Self::checked_sandbox_id("sandbox.readFile", sandbox_id)?;
let body = self
.execute_op(
sandbox_id,
"sandbox.readFile",
json!({ "v": AGENT_PROTOCOL_VERSION, "op": "readFile", "path": path }),
)
.await?;
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
struct ReadFile {
contents_base64: String,
}
let read: ReadFile = serde_json::from_slice(&body).map_err(|_| {
AlienError::new(ErrorData::SandboxCommandFailed {
failure: "agentRefused".to_string(),
reason: format!("sandbox.readFile was refused: {}", truncated(&body)),
})
})?;
BASE64
.decode(read.contents_base64.as_bytes())
.map_err(|error| {
AlienError::new(ErrorData::UnexpectedResponseFormat {
provider: "gcp-agent-platform".to_string(),
binding_name: "sandbox.readFile".to_string(),
field: "contentsBase64".to_string(),
response_json: format!("the agent returned data that is not base64: {error}"),
})
})
}
async fn write_files(&self, sandbox_id: &str, files: BTreeMap<String, Vec<u8>>) -> Result<()> {
Self::checked_sandbox_id("sandbox.writeFiles", sandbox_id)?;
for (path, contents) in files {
let body = self
.execute_op(
sandbox_id,
"sandbox.writeFiles",
json!({
"v": AGENT_PROTOCOL_VERSION,
"op": "writeFile",
"path": path,
"contentsBase64": BASE64.encode(&contents),
}),
)
.await?;
confirm_empty_ok("sandbox.writeFiles", &body)?;
}
Ok(())
}
async fn preview(&self, _sandbox_id: &str, _port: u16) -> Result<PreviewCapability> {
Err(self.unsupported(
"preview",
"Agent Platform mints no port-scoped ingress capability; the only ingress is :execute",
))
}
async fn pause(&self, sandbox_id: &str) -> Result<()> {
Self::checked_sandbox_id("sandbox.pause", sandbox_id)?;
Err(self.pause_resume_unsupported())
}
async fn resume(&self, sandbox_id: &str) -> Result<()> {
Self::checked_sandbox_id("sandbox.resume", sandbox_id)?;
Err(self.pause_resume_unsupported())
}
async fn snapshot(&self, sandbox_id: &str) -> Result<String> {
Self::checked_sandbox_id("sandbox.snapshot", sandbox_id)?;
let display_name = format!("snap-{}", uuid::Uuid::new_v4().simple());
let started = self
.client
.snapshot(&self.engine, sandbox_id, &display_name)
.await
.context(ErrorData::SandboxCommandFailed {
failure: "snapshotFailed".to_string(),
reason: format!("sandbox.snapshot: sandbox '{sandbox_id}' could not be captured"),
})?;
let snapshot: SandboxSnapshot =
serde_json::from_value(self.await_operation("sandbox.snapshot", started).await?)
.map_err(|error| {
AlienError::new(ErrorData::UnexpectedResponseFormat {
provider: "gcp-agent-platform".to_string(),
binding_name: "sandbox.snapshot".to_string(),
field: "response".to_string(),
response_json: format!(
"the snapshot operation resolved to a non-snapshot: {error}"
),
})
})?;
snapshot.name.ok_or_else(|| {
AlienError::new(ErrorData::UnexpectedResponseFormat {
provider: "gcp-agent-platform".to_string(),
binding_name: "sandbox.snapshot".to_string(),
field: "name".to_string(),
response_json: "the snapshot completed without a resource name".to_string(),
})
})
}
async fn terminate(&self, sandbox_id: &str) -> Result<()> {
Self::checked_sandbox_id(TERMINATE, sandbox_id)?;
if let Err(error) = self.client.delete_sandbox(&self.engine, sandbox_id).await {
if !is_not_found(&error) {
return Err(error.context(ErrorData::SandboxUnreachable {
operation: TERMINATE.to_string(),
reason: format!("the delete of sandbox '{sandbox_id}' was not accepted"),
}));
}
return Ok(());
}
for _ in 0..TERMINATE_POLL_ATTEMPTS {
match self.client.get_sandbox(&self.engine, sandbox_id).await {
Err(error) if is_not_found(&error) => return Ok(()),
Err(error) => {
warn!(sandbox = %sandbox_id, %error, "could not confirm a sandbox is gone")
}
Ok(_) => {}
}
tokio::time::sleep(TERMINATE_POLL_INTERVAL).await;
}
Err(AlienError::new(ErrorData::SandboxUnreachable {
operation: TERMINATE.to_string(),
reason: format!(
"deletion of '{sandbox_id}' was accepted but the sandbox was still present after \
{}s; it may still be running",
TERMINATE_POLL_ATTEMPTS as u64 * TERMINATE_POLL_INTERVAL.as_secs()
),
}))
}
}
async fn job_poll_step(mut state: JobPollState) -> Option<(Result<CommandOutput>, JobPollState)> {
loop {
if let Some(item) = state.pending.pop_front() {
return Some((item, state));
}
if state.finished {
return None;
}
if tokio::time::Instant::now() >= state.deadline_at {
let cancelled = state
.client
.execute(
&state.engine,
&state.sandbox_id,
&cancel_body(&state.job_id),
)
.await;
let confirmed = cancelled.as_ref().is_ok_and(|body| cancel_confirmed(body));
state.pending.push_back(Err(match cancelled {
Ok(_) if confirmed => AlienError::new(ErrorData::SandboxCommandFailed {
failure: "timeoutExceeded".to_string(),
reason: "the command's deadline elapsed before its job reported an outcome"
.to_string(),
}),
Ok(_) => AlienError::new(ErrorData::SandboxOutcomeUnknown {
operation: RUN_COMMAND.to_string(),
reason: "the command's deadline elapsed and its job did not confirm the cancel"
.to_string(),
}),
Err(error) => error.context(ErrorData::SandboxOutcomeUnknown {
operation: RUN_COMMAND.to_string(),
reason: "the command's deadline elapsed and its job could not be cancelled"
.to_string(),
}),
}));
state.finished = true;
state.stopped = confirmed;
continue;
}
let poll = match poll_once(
state.client.as_ref(),
&state.engine,
&state.sandbox_id,
&state.job_id,
state.since_seq,
)
.await
{
Ok(poll) => poll,
Err(error) => {
state
.pending
.push_back(Err(error.context(ErrorData::SandboxOutcomeUnknown {
operation: RUN_COMMAND.to_string(),
reason: "the job is no longer watched".to_string(),
})));
state.finished = true;
continue;
}
};
for frame in poll.frames {
state.since_seq = state.since_seq.max(frame.seq());
let output = frame.into_output();
let failed = output.is_err();
state.pending.push_back(output);
if failed {
state.finished = true;
break;
}
}
if !poll.running {
let terminal = match poll.error {
Some(error) => Err(AlienError::new(ErrorData::SandboxCommandFailed {
failure: error.code,
reason: error.message,
})),
None => match poll.exit_code {
Some(code) => Ok(CommandOutput::Exit {
code,
truncated: poll.truncated.unwrap_or(false),
}),
None => Err(AlienError::new(ErrorData::SandboxOutcomeUnknown {
operation: RUN_COMMAND.to_string(),
reason: "the job finished without reporting an exit code".to_string(),
})),
},
};
state.pending.push_back(terminal);
state.finished = true;
state.stopped = true;
continue;
}
if state.pending.is_empty() {
tokio::time::sleep(JOB_POLL_INTERVAL).await;
}
}
}
async fn poll_once(
client: &dyn AgentPlatformApi,
engine: &str,
sandbox_id: &str,
job_id: &str,
since_seq: Option<u64>,
) -> Result<JobPollReply> {
let body = client
.execute(engine, sandbox_id, &poll_body(job_id, since_seq))
.await
.map_err(|error| unanswered_job(JOB_POLL, error))?;
serde_json::from_slice(&body).map_err(|_| {
AlienError::new(ErrorData::UnexpectedResponseFormat {
provider: "gcp-agent-platform".to_string(),
binding_name: JOB_POLL.to_string(),
field: "jobPoll".to_string(),
response_json: truncated(&body),
})
})
}
fn unanswered_job(
operation: &str,
error: AlienError<AgentPlatformErrorData>,
) -> AlienError<ErrorData> {
if is_not_found(&error) {
return error.context(ErrorData::SandboxCommandFailed {
failure: "sandboxGone".to_string(),
reason: format!("{operation}: the sandbox does not exist"),
});
}
error.context(ErrorData::SandboxUnreachable {
operation: operation.to_string(),
reason: "the sandbox did not complete the call".to_string(),
})
}
fn cancel_confirmed(body: &[u8]) -> bool {
serde_json::from_slice::<serde_json::Value>(body)
.is_ok_and(|value| value.as_object().is_some_and(serde_json::Map::is_empty))
}
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
struct JobStartReply {
job_id: String,
}
struct JobPollState {
client: Arc<dyn AgentPlatformApi>,
engine: String,
sandbox_id: String,
job_id: String,
since_seq: Option<u64>,
pending: VecDeque<Result<CommandOutput>>,
finished: bool,
stopped: bool,
deadline_at: tokio::time::Instant,
}
impl Drop for JobPollState {
fn drop(&mut self) {
if self.stopped {
return;
}
let Ok(runtime) = tokio::runtime::Handle::try_current() else {
warn!(
sandbox = %self.sandbox_id,
job = %self.job_id,
"a job's stream was dropped outside a runtime, so the job runs to its timeout"
);
return;
};
let client = self.client.clone();
let engine = self.engine.clone();
let sandbox_id = self.sandbox_id.clone();
let job_id = self.job_id.clone();
runtime.spawn(async move {
match tokio::time::timeout(
AGENT_PROBE_BUDGET,
client.execute(&engine, &sandbox_id, &cancel_body(&job_id)),
)
.await
{
Ok(Ok(body)) if cancel_confirmed(&body) => {}
Ok(Ok(body)) => warn!(
sandbox = %sandbox_id,
job = %job_id,
reply = %truncated(&body),
"a dropped job's cancel was refused, so the job runs to its timeout"
),
Ok(Err(error)) => warn!(
sandbox = %sandbox_id,
job = %job_id,
%error,
"a dropped job's cancel did not land, so the job may run to its timeout"
),
Err(_) => warn!(
sandbox = %sandbox_id,
job = %job_id,
budget_secs = AGENT_PROBE_BUDGET.as_secs(),
"a dropped job's cancel went unanswered, so the job may run to its timeout"
),
}
});
}
}
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
struct JobPollReply {
running: bool,
#[serde(default)]
frames: Vec<WireFrame>,
#[serde(default)]
exit_code: Option<i32>,
#[serde(default)]
truncated: Option<bool>,
#[serde(default)]
error: Option<JobErrorReply>,
}
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
struct JobErrorReply {
code: String,
message: String,
}
#[derive(Deserialize)]
#[serde(rename_all = "camelCase", tag = "t")]
enum WireFrame {
Stdout {
seq: u64,
data: String,
},
Stderr {
seq: u64,
data: String,
},
Exit {
code: i32,
#[serde(default)]
truncated: bool,
},
Error {
code: String,
message: String,
},
}
impl WireFrame {
fn is_terminal(&self) -> bool {
matches!(self, Self::Exit { .. } | Self::Error { .. })
}
fn seq(&self) -> Option<u64> {
match self {
Self::Stdout { seq, .. } | Self::Stderr { seq, .. } => Some(*seq),
_ => None,
}
}
fn into_output(self) -> Result<CommandOutput> {
match self {
Self::Stdout { seq, data } => Ok(CommandOutput::Stdout {
seq,
data: decode_frame_data(&data)?,
}),
Self::Stderr { seq, data } => Ok(CommandOutput::Stderr {
seq,
data: decode_frame_data(&data)?,
}),
Self::Exit { code, truncated } => Ok(CommandOutput::Exit { code, truncated }),
Self::Error { code, message } => {
Err(AlienError::new(ErrorData::SandboxCommandFailed {
failure: code,
reason: message,
}))
}
}
}
}
fn decode_frame_data(data: &str) -> Result<Vec<u8>> {
BASE64.decode(data).map_err(|error| {
AlienError::new(ErrorData::UnexpectedResponseFormat {
provider: "gcp-agent-platform".to_string(),
binding_name: RUN_COMMAND.to_string(),
field: "data".to_string(),
response_json: format!("an output frame's data is not base64: {error}"),
})
.context(ErrorData::SandboxOutcomeUnknown {
operation: RUN_COMMAND.to_string(),
reason: "an output frame did not decode".to_string(),
})
})
}
fn parse_exec_frames(body: &[u8]) -> Result<Vec<Result<CommandOutput>>> {
let mut frames = Vec::new();
let mut saw_any = false;
let mut saw_terminal = false;
for line in body.split(|byte| *byte == b'\n') {
if line.is_empty() {
continue;
}
match serde_json::from_slice::<WireFrame>(line) {
Ok(frame) => {
saw_any = true;
saw_terminal |= frame.is_terminal();
let output = frame.into_output();
let failed = output.is_err();
frames.push(output);
if failed {
saw_terminal = true;
break;
}
}
Err(error) => {
if !saw_any {
return Err(AlienError::new(ErrorData::SandboxCommandFailed {
failure: "agentRefused".to_string(),
reason: format!("run_command was refused: {}", truncated(body)),
}));
}
frames.push(Err(AlienError::new(ErrorData::UnexpectedResponseFormat {
provider: "gcp-agent-platform".to_string(),
binding_name: RUN_COMMAND.to_string(),
field: "frame".to_string(),
response_json: format!("an output frame did not parse: {error}"),
})
.context(ErrorData::SandboxOutcomeUnknown {
operation: RUN_COMMAND.to_string(),
reason: "an output frame did not parse".to_string(),
})));
saw_terminal = true;
break;
}
}
}
if !saw_any {
return Err(AlienError::new(ErrorData::SandboxCommandFailed {
failure: "agentRefused".to_string(),
reason: "run_command returned an empty body".to_string(),
}));
}
if !saw_terminal {
frames.push(Err(AlienError::new(ErrorData::SandboxOutcomeUnknown {
operation: RUN_COMMAND.to_string(),
reason: "the command's output ended without a terminal frame".to_string(),
})));
}
Ok(frames)
}
fn exec_envelope(op: &str, _sandbox_id: &str, request: &RunCommandRequest) -> serde_json::Value {
json!({
"v": AGENT_PROTOCOL_VERSION,
"op": op,
"command": request.argv(),
"timeoutMs": timeout_millis(request.timeout),
"cwd": request.cwd,
"env": request.env,
})
}
fn poll_body(job_id: &str, since_seq: Option<u64>) -> Vec<u8> {
serde_json::to_vec(&json!({
"v": AGENT_PROTOCOL_VERSION,
"op": "jobPoll",
"jobId": job_id,
"sinceSeq": since_seq,
}))
.unwrap_or_default()
}
fn cancel_body(job_id: &str) -> Vec<u8> {
serde_json::to_vec(&json!({
"v": AGENT_PROTOCOL_VERSION,
"op": "jobCancel",
"jobId": job_id,
}))
.unwrap_or_default()
}
fn timeout_millis(timeout: Duration) -> u64 {
u64::try_from(timeout.as_millis()).unwrap_or(u64::MAX)
}
fn confirm_empty_ok(operation: &str, body: &[u8]) -> Result<()> {
if body.iter().all(|byte| byte.is_ascii_whitespace()) {
return Ok(());
}
Err(AlienError::new(ErrorData::SandboxCommandFailed {
failure: "agentRefused".to_string(),
reason: format!("{operation} was refused: {}", truncated(body)),
}))
}
fn sandbox_segment(name: &str) -> Option<&str> {
let segment = name.rsplit('/').next()?;
is_addressable_id(segment).then_some(segment)
}
fn is_addressable_id(id: &str) -> bool {
!id.is_empty()
&& id.len() <= MAX_SANDBOX_ID
&& id
.chars()
.all(|c| c.is_ascii_alphanumeric() || c == '-' || c == '_')
}
fn generation_from_boot_id(boot_id: &str) -> u64 {
const FNV_OFFSET_BASIS: u64 = 0xcbf2_9ce4_8422_2325;
const FNV_PRIME: u64 = 0x0000_0100_0000_01b3;
let mut hash = FNV_OFFSET_BASIS;
for byte in boot_id.as_bytes() {
hash ^= u64::from(*byte);
hash = hash.wrapping_mul(FNV_PRIME);
}
hash | 1
}
fn sandbox_state(operation: &str, state: Option<&str>) -> Result<SandboxState> {
match state {
Some("STATE_RUNNING") => Ok(SandboxState::Running),
Some("STATE_CREATING" | "STATE_PENDING" | "STATE_RESUMING") => Ok(SandboxState::Starting),
Some("STATE_PAUSED" | "STATE_PAUSING" | "STATE_SUSPENDED") => Ok(SandboxState::Paused),
Some("STATE_STOPPED" | "STATE_FAILED" | "STATE_DELETING" | "STATE_DELETED") => {
Ok(SandboxState::Terminated)
}
other => Err(AlienError::new(ErrorData::UnexpectedResponseFormat {
provider: "gcp-agent-platform".to_string(),
binding_name: operation.to_string(),
field: "state".to_string(),
response_json: other
.map_or_else(|| "absent".to_string(), |state| format!("\"{state}\"")),
})),
}
}
fn finish_operation(operation: &str, name: &str, op: Operation) -> Result<serde_json::Value> {
match op.result {
Some(OperationResult::Response { response }) => Ok(response),
Some(OperationResult::Error { error }) => {
Err(AlienError::new(ErrorData::SandboxCommandFailed {
failure: "operationFailed".to_string(),
reason: format!(
"{operation}: operation '{name}' failed (grpc {}): {}",
error.code, error.message
),
}))
}
None => Err(AlienError::new(ErrorData::UnexpectedResponseFormat {
provider: "gcp-agent-platform".to_string(),
binding_name: operation.to_string(),
field: "response".to_string(),
response_json: format!("operation '{name}' reported done without a result"),
})),
}
}
fn agent_answer(error: &AlienError<AgentPlatformErrorData>) -> Option<String> {
const RELAYED: &str = "Error Details: ";
let message = super::refusal::captured_refusal(error)?;
let details = message.split_once(RELAYED)?.1.trim();
let code = details.split(':').next()?;
let is_agent_code =
!code.is_empty() && code.chars().all(|ch| ch.is_ascii_uppercase() || ch == '_');
is_agent_code.then(|| truncated(details.as_bytes()))
}
fn is_not_found(error: &AlienError<AgentPlatformErrorData>) -> bool {
const NOT_FOUND: &str = "REMOTE_RESOURCE_NOT_FOUND";
if error.code == NOT_FOUND {
return true;
}
let mut node = error.source.as_deref();
while let Some(current) = node {
if current.code == NOT_FOUND {
return true;
}
node = current.source.as_deref();
}
false
}
fn truncated(body: &[u8]) -> String {
const LIMIT: usize = 200;
let text = String::from_utf8_lossy(body);
let text = text.trim();
if text.len() <= LIMIT {
return text.to_string();
}
let end = (0..=LIMIT)
.rev()
.find(|at| text.is_char_boundary(*at))
.unwrap_or(0);
format!("{}…", &text[..end])
}
#[cfg(test)]
#[path = "gcp_agent_platform_tests.rs"]
mod tests;