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, CreateSessionRequest, PreviewCapability, RunCommandRequest, Sandbox,
SandboxSession, SandboxSessionState,
};
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_DEADLINE: Duration = Duration::from_secs(30);
const MAX_SESSION_ID: usize = 63;
const SESSION_READY_ATTEMPTS: u32 = 150;
const SESSION_READY_INTERVAL: Duration = Duration::from_secs(2);
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 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,
session_ttl_seconds: Option<u32>,
}
impl GcpAgentPlatformSandbox {
pub fn new(
client: Arc<dyn AgentPlatformApi>,
engine: String,
template: String,
session_ttl_seconds: Option<u32>,
) -> Self {
let engine = engine.rsplit('/').next().unwrap_or(&engine).to_string();
Self {
client,
engine,
template,
session_ttl_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 checked_session_id(operation: &str, session_id: &str) -> Result<()> {
if is_addressable_id(session_id) {
return Ok(());
}
Err(AlienError::new(ErrorData::InvalidInput {
operation_context: operation.to_string(),
details: format!(
"session id '{session_id}' must be a single segment of letters, digits, '-' and \
'_', at most {MAX_SESSION_ID} characters"
),
field_name: Some("sessionId".to_string()),
}))
}
async fn read_sandbox(
&self,
operation: &str,
session_id: &str,
) -> Result<Option<SandboxEnvironment>> {
match self.client.get_sandbox(&self.engine, session_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,
session_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, session_id, &body)
.await
.map_err(|error| Self::execute_failed(operation, error))
}
fn execute_failed(
operation: &str,
error: AlienError<AgentPlatformErrorData>,
) -> AlienError<ErrorData> {
if is_not_found(&error) {
return error.context(ErrorData::SandboxCommandFailed {
failure: "sessionGone".to_string(),
reason: format!("{operation}: the session does not exist"),
});
}
error.context(ErrorData::SandboxCommandFailed {
failure: "executeFailed".to_string(),
reason: format!("{operation} could not be completed against the session"),
})
}
async fn probe_agent(&self, operation: &str, session_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,
session_id,
&serde_json::to_vec(&json!({ "v": AGENT_PROTOCOL_VERSION, "op": "health" }))
.unwrap_or_default(),
),
)
.await
.map_err(|_| {
unreachable(format!(
"the session'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 session'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 session'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 session's agent speaks protocol {} where this provider speaks {}",
health.protocol_version, AGENT_PROTOCOL_VERSION
)));
}
if health.boot_id.is_empty() {
return Err(unreachable(
"the session'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,
session_id: &str,
reason: AlienError<ErrorData>,
) -> AlienError<ErrorData> {
let Err(error) = self.client.delete_sandbox(&self.engine, session_id).await else {
return reason;
};
warn!(
session = %session_id,
%error,
"could not delete a sandbox that was never handed to its caller"
);
reason.context(ErrorData::SandboxCommandFailed {
failure: "sandboxLeftBehind".to_string(),
reason: format!(
"session '{session_id}' was not handed to its caller and could not be deleted, so \
it is still running"
),
})
}
async fn settle(&self, session_id: &str) -> Result<u64> {
for _ in 0..SESSION_READY_ATTEMPTS {
let Some(sandbox) = self.read_sandbox(CREATE, session_id).await? else {
return Err(AlienError::new(ErrorData::SandboxCommandFailed {
failure: "sessionGone".to_string(),
reason: format!("session '{session_id}' disappeared while it was coming up"),
}));
};
match session_state(CREATE, sandbox.state.as_deref())? {
SandboxSessionState::Running => {
return self.probe_agent(CREATE, session_id).await;
}
SandboxSessionState::Terminated => {
return Err(AlienError::new(ErrorData::SandboxCommandFailed {
failure: "sessionTerminated".to_string(),
reason: format!(
"session '{session_id}' reached a terminal state while starting"
),
}));
}
SandboxSessionState::Starting | SandboxSessionState::Suspended => {}
}
tokio::time::sleep(SESSION_READY_INTERVAL).await;
}
Err(AlienError::new(ErrorData::SandboxUnreachable {
operation: CREATE.to_string(),
reason: format!(
"session '{session_id}' was not running after {}s",
SESSION_READY_ATTEMPTS as u64 * SESSION_READY_INTERVAL.as_secs()
),
}))
}
async fn run_synchronous(
&self,
session_id: &str,
request: &RunCommandRequest,
) -> Result<BoxStream<'static, Result<CommandOutput>>> {
let envelope = exec_envelope("exec", session_id, request);
let body = self.execute_op(session_id, RUN_COMMAND, envelope).await?;
let frames = parse_exec_frames(&body)?;
Ok(Box::pin(stream::iter(frames)))
}
async fn run_detached(
&self,
session_id: &str,
request: &RunCommandRequest,
) -> Result<BoxStream<'static, Result<CommandOutput>>> {
let envelope = exec_envelope("jobStart", session_id, request);
let body = self.execute_op(session_id, RUN_COMMAND, envelope).await?;
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
struct JobStart {
job_id: String,
}
let started: JobStart = serde_json::from_slice(&body).map_err(|_| {
AlienError::new(ErrorData::UnexpectedResponseFormat {
provider: "gcp-agent-platform".to_string(),
binding_name: RUN_COMMAND.to_string(),
field: "jobId".to_string(),
response_json: truncated(&body),
})
})?;
let state = JobPollState {
client: self.client.clone(),
engine: self.engine.clone(),
session_id: session_id.to_string(),
job_id: started.job_id,
since_seq: None,
pending: VecDeque::new(),
finished: false,
deadline_at: tokio::time::Instant::now() + request.deadline + JOB_POLL_GRACE,
};
Ok(Box::pin(stream::unfold(state, job_poll_step)))
}
}
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: CreateSessionRequest) -> Result<SandboxSession> {
if !request.env.is_empty() {
return Err(AlienError::new(ErrorData::InvalidInput {
operation_context: CREATE.to_string(),
details: "Agent Platform carries no per-session environment; pass variables on \
each command instead"
.to_string(),
field_name: Some("env".to_string()),
}));
}
let started = self
.client
.create_sandbox(
&self.engine,
SandboxCreateRequest {
display_name: request.session_id.clone(),
sandbox_environment_template: Some(self.template.clone()),
sandbox_environment_snapshot: None,
ttl: self
.session_ttl_seconds
.map(|seconds| format!("{seconds}s")),
},
)
.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(session_id) = created.name.as_deref().and_then(session_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 session_id = session_id.to_string();
match self.settle(&session_id).await {
Ok(generation) => Ok(SandboxSession {
session_id,
state: SandboxSessionState::Running,
generation,
}),
Err(error) => Err(self.discard(&session_id, error).await),
}
}
async fn get(&self, session_id: &str) -> Result<Option<SandboxSession>> {
Self::checked_session_id(GET, session_id)?;
let Some(sandbox) = self.read_sandbox(GET, session_id).await? else {
return Ok(None);
};
let state = session_state(GET, sandbox.state.as_deref())?;
let generation = if state == SandboxSessionState::Running {
self.probe_agent(GET, session_id).await?
} else {
NO_GENERATION
};
Ok(Some(SandboxSession {
session_id: session_id.to_string(),
state,
generation,
}))
}
async fn get_or_create(&self, request: CreateSessionRequest) -> Result<SandboxSession> {
if let Some(id) = request.session_id.as_deref() {
match self.get(id).await {
Ok(Some(session)) if session.state == SandboxSessionState::Running => {
return Ok(session)
}
Ok(Some(session)) if session.state == SandboxSessionState::Suspended => {
if self.resume(id).await.is_ok() {
match self.get(id).await {
Ok(Some(woken)) if woken.state == SandboxSessionState::Running => {
return Ok(woken)
}
_ => {
if let Err(error) = self.suspend(id).await {
return Err(error.context(ErrorData::SandboxCommandFailed {
failure: "resumeRollbackFailed".to_string(),
reason: format!(
"{GET_OR_CREATE}: woke session '{id}' but could not \
confirm it healthy or put it back to sleep"
),
}));
}
}
}
}
}
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 session '{id}' failed"),
}))
}
}
}
self.create(request).await
}
async fn list(&self) -> Result<Vec<SandboxSession>> {
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 session_id = sandbox.name.as_deref().and_then(session_segment)?;
let state = session_state("sandbox.list", sandbox.state.as_deref()).ok()?;
Some(SandboxSession {
session_id: session_id.to_string(),
state,
generation: NO_GENERATION,
})
})
.collect())
}
async fn run_command(
&self,
session_id: &str,
request: RunCommandRequest,
) -> Result<BoxStream<'static, Result<CommandOutput>>> {
Self::checked_session_id(RUN_COMMAND, session_id)?;
if request.command.is_empty() {
return Err(AlienError::new(ErrorData::InvalidInput {
operation_context: RUN_COMMAND.to_string(),
details: "a command must name a program to run".to_string(),
field_name: Some("command".to_string()),
}));
}
if deadline_millis(request.deadline) == 0 {
return Err(AlienError::new(ErrorData::SandboxCommandFailed {
failure: "invalidRequest".to_string(),
reason: "a command must carry a deadline of at least one millisecond".to_string(),
}));
}
if request.deadline <= MAX_SYNCHRONOUS_DEADLINE {
self.run_synchronous(session_id, &request).await
} else {
self.run_detached(session_id, &request).await
}
}
async fn read_file(&self, session_id: &str, path: &str) -> Result<Vec<u8>> {
Self::checked_session_id("sandbox.readFile", session_id)?;
let body = self
.execute_op(
session_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, session_id: &str, files: BTreeMap<String, Vec<u8>>) -> Result<()> {
Self::checked_session_id("sandbox.writeFiles", session_id)?;
for (path, contents) in files {
let body = self
.execute_op(
session_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 mkdir(&self, session_id: &str, path: &str) -> Result<()> {
Self::checked_session_id("sandbox.mkdir", session_id)?;
let body = self
.execute_op(
session_id,
"sandbox.mkdir",
json!({ "v": AGENT_PROTOCOL_VERSION, "op": "mkdir", "path": path }),
)
.await?;
confirm_empty_ok("sandbox.mkdir", &body)
}
async fn preview(&self, _session_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 suspend(&self, session_id: &str) -> Result<()> {
Self::checked_session_id("sandbox.suspend", session_id)?;
let started = self.client.pause(&self.engine, session_id).await.context(
ErrorData::SandboxCommandFailed {
failure: "suspendFailed".to_string(),
reason: format!("sandbox.suspend: session '{session_id}' could not be paused"),
},
)?;
self.await_operation("sandbox.suspend", started).await?;
Ok(())
}
async fn resume(&self, session_id: &str) -> Result<()> {
Self::checked_session_id("sandbox.resume", session_id)?;
let started = self.client.resume(&self.engine, session_id).await.context(
ErrorData::SandboxCommandFailed {
failure: "resumeFailed".to_string(),
reason: format!("sandbox.resume: session '{session_id}' could not be resumed"),
},
)?;
self.await_operation("sandbox.resume", started).await?;
Ok(())
}
async fn snapshot(&self, session_id: &str) -> Result<String> {
Self::checked_session_id("sandbox.snapshot", session_id)?;
let display_name = format!("snap-{}", uuid::Uuid::new_v4().simple());
let started = self
.client
.snapshot(&self.engine, session_id, &display_name)
.await
.context(ErrorData::SandboxCommandFailed {
failure: "snapshotFailed".to_string(),
reason: format!("sandbox.snapshot: session '{session_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, session_id: &str) -> Result<()> {
Self::checked_session_id(TERMINATE, session_id)?;
self.client
.delete_sandbox(&self.engine, session_id)
.await
.context(ErrorData::SandboxUnreachable {
operation: TERMINATE.to_string(),
reason: format!("the delete of session '{session_id}' was not accepted"),
})?;
for _ in 0..TERMINATE_POLL_ATTEMPTS {
match self.client.get_sandbox(&self.engine, session_id).await {
Err(error) if is_not_found(&error) => return Ok(()),
Err(error) => {
warn!(session = %session_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 '{session_id}' was accepted but the session 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 _ = state
.client
.execute(
&state.engine,
&state.session_id,
&cancel_body(&state.job_id),
)
.await;
state
.pending
.push_back(Err(AlienError::new(ErrorData::SandboxCommandFailed {
failure: "deadlineExceeded".to_string(),
reason: "the command's deadline elapsed before its job reported an outcome"
.to_string(),
})));
state.finished = true;
continue;
}
let body = match state
.client
.execute(
&state.engine,
&state.session_id,
&poll_body(&state.job_id, state.since_seq),
)
.await
{
Ok(body) => body,
Err(error) => {
state
.pending
.push_back(Err(GcpAgentPlatformSandbox::execute_failed(
RUN_COMMAND,
error,
)));
state.finished = true;
continue;
}
};
let poll: JobPoll = match serde_json::from_slice(&body) {
Ok(poll) => poll,
Err(_) => {
state.pending.push_back(Err(AlienError::new(
ErrorData::UnexpectedResponseFormat {
provider: "gcp-agent-platform".to_string(),
binding_name: RUN_COMMAND.to_string(),
field: "jobPoll".to_string(),
response_json: truncated(&body),
},
)));
state.finished = true;
continue;
}
};
for frame in poll.frames {
state.since_seq = state.since_seq.max(frame.seq());
state.pending.push_back(frame.into_output());
}
if !poll.running {
let terminal = match poll.error {
Some(error) => Err(AlienError::new(ErrorData::SandboxCommandFailed {
failure: error.code,
reason: error.message,
})),
None => Ok(CommandOutput::Exit {
code: poll.exit_code.unwrap_or(-1),
truncated: poll.truncated.unwrap_or(false),
}),
};
state.pending.push_back(terminal);
state.finished = true;
continue;
}
if state.pending.is_empty() {
tokio::time::sleep(JOB_POLL_INTERVAL).await;
}
}
}
struct JobPollState {
client: Arc<dyn AgentPlatformApi>,
engine: String,
session_id: String,
job_id: String,
since_seq: Option<u64>,
pending: VecDeque<Result<CommandOutput>>,
finished: bool,
deadline_at: tokio::time::Instant,
}
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
struct JobPoll {
running: bool,
#[serde(default)]
frames: Vec<WireFrame>,
#[serde(default)]
exit_code: Option<i32>,
#[serde(default)]
truncated: Option<bool>,
#[serde(default)]
error: Option<JobError>,
}
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
struct JobError {
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}"),
})
})
}
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();
frames.push(frame.into_output());
}
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}"),
})));
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::SandboxCommandFailed {
failure: "outcomeUnknown".to_string(),
reason: "the command's output ended without a terminal frame, so whether it finished \
is unknown"
.to_string(),
})));
}
Ok(frames)
}
fn exec_envelope(op: &str, _session_id: &str, request: &RunCommandRequest) -> serde_json::Value {
json!({
"v": AGENT_PROTOCOL_VERSION,
"op": op,
"command": request.command,
"deadlineMs": deadline_millis(request.deadline),
"workingDirectory": request.working_directory,
"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 deadline_millis(deadline: Duration) -> u64 {
u64::try_from(deadline.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 session_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_SESSION_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 session_state(operation: &str, state: Option<&str>) -> Result<SandboxSessionState> {
match state {
Some("STATE_RUNNING") => Ok(SandboxSessionState::Running),
Some("STATE_CREATING" | "STATE_PENDING" | "STATE_RESUMING") => {
Ok(SandboxSessionState::Starting)
}
Some("STATE_PAUSED" | "STATE_PAUSING" | "STATE_SUSPENDED") => {
Ok(SandboxSessionState::Suspended)
}
Some("STATE_STOPPED" | "STATE_FAILED" | "STATE_DELETING" | "STATE_DELETED") => {
Ok(SandboxSessionState::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 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;