use super::*;
use alien_gcp_clients::gcp::agent_platform::{
MockAgentPlatformApi, ReasoningEngine, SandboxEnvironmentTemplate,
};
use futures::StreamExt;
use std::sync::atomic::{AtomicUsize, Ordering};
type ClientResult<T> = alien_error::Result<T, AgentPlatformErrorData>;
const ENGINE_FULL: &str = "projects/p/locations/us-central1/reasoningEngines/eng1";
const TEMPLATE: &str = "projects/p/locations/us-central1/sandboxTemplates/agent";
fn provider(client: MockAgentPlatformApi) -> GcpAgentPlatformSandbox {
provider_from(Arc::new(client))
}
fn provider_from(client: Arc<dyn AgentPlatformApi>) -> GcpAgentPlatformSandbox {
GcpAgentPlatformSandbox::new(
client,
ENGINE_FULL.to_string(),
TEMPLATE.to_string(),
Some(3600),
)
}
fn sandbox_name(id: &str) -> String {
format!("{ENGINE_FULL}/sandboxEnvironments/{id}")
}
fn sandbox_in_state(id: &str, state: &str) -> SandboxEnvironment {
SandboxEnvironment {
name: Some(sandbox_name(id)),
display_name: None,
state: Some(state.to_string()),
sandbox_environment_template: None,
expire_time: None,
connection_info: None,
extra: Default::default(),
}
}
fn done_op(value: serde_json::Value) -> Operation {
Operation {
name: Some("projects/p/locations/us-central1/operations/op1".to_string()),
metadata: None,
done: Some(true),
result: Some(OperationResult::Response { response: value }),
}
}
fn op_of(input: &[u8]) -> String {
serde_json::from_slice::<serde_json::Value>(input)
.ok()
.and_then(|value| {
value
.get("op")
.and_then(|op| op.as_str())
.map(str::to_string)
})
.unwrap_or_default()
}
fn ndjson(lines: &[serde_json::Value]) -> Vec<u8> {
let mut body = Vec::new();
for line in lines {
body.extend_from_slice(
serde_json::to_string(line)
.expect("frame serializes")
.as_bytes(),
);
body.push(b'\n');
}
body
}
fn health_reply() -> Vec<u8> {
health_reply_with_boot("11111111-1111-1111-1111-111111111111")
}
fn health_reply_with_boot(boot_id: &str) -> Vec<u8> {
serde_json::to_vec(&serde_json::json!({ "protocolVersion": 1, "bootId": boot_id }))
.expect("health serializes")
}
fn stdout_frame(seq: u64, data: &[u8]) -> serde_json::Value {
serde_json::json!({ "t": "stdout", "seq": seq, "data": BASE64.encode(data) })
}
fn exit_frame(code: i32) -> serde_json::Value {
serde_json::json!({ "t": "exit", "code": code, "truncated": false })
}
fn not_found() -> AlienError<AgentPlatformErrorData> {
AlienError::new(alien_client_core::ErrorData::RemoteResourceNotFound {
resource_type: "SandboxEnvironment".to_string(),
resource_name: "s1".to_string(),
})
.context(AgentPlatformErrorData::RequestFailed {
operation: "get sandbox".to_string(),
message: "s1".to_string(),
})
}
fn execute_refused() -> AlienError<AgentPlatformErrorData> {
AlienError::new(AgentPlatformErrorData::ExecuteFailed {
sandbox: "s1".to_string(),
message: "the API rejected the request".to_string(),
})
}
#[tokio::test]
async fn create_awaits_running_probes_the_agent_and_pins_its_arguments() {
let mut client = MockAgentPlatformApi::new();
client
.expect_create_sandbox()
.withf(|engine, request| {
engine == "eng1"
&& request.sandbox_environment_template.as_deref() == Some(TEMPLATE)
&& request.ttl.as_deref() == Some("3600s")
})
.times(1)
.returning(|_, _| Ok(done_op(serde_json::json!({ "name": sandbox_name("s1") }))));
client
.expect_get_sandbox()
.withf(|engine, sandbox| engine == "eng1" && sandbox == "s1")
.returning(|_, id| Ok(sandbox_in_state(id, "STATE_RUNNING")));
client
.expect_execute()
.withf(|_, sandbox, input| sandbox == "s1" && op_of(input) == "health")
.returning(|_, _, _| Ok(health_reply()));
let sandbox = provider(client)
.create(CreateSandboxRequest::default())
.await
.expect("create succeeds");
assert_eq!(sandbox.sandbox_id, "s1");
assert_eq!(sandbox.state, SandboxState::Running);
}
#[tokio::test(start_paused = true)]
async fn create_waits_for_the_agent_and_deletes_the_sandbox_when_it_never_answers() {
fn bad_gateway() -> AlienError<AgentPlatformErrorData> {
AlienError::new(AgentPlatformErrorData::ExecuteFailed {
sandbox: "s1".to_string(),
message: "Bad Gateway: Unable to reach the sandbox environment.".to_string(),
})
}
let client_with = |answer_after: Option<usize>, deletes: usize| {
let mut client = MockAgentPlatformApi::new();
client
.expect_create_sandbox()
.returning(|_, _| Ok(done_op(serde_json::json!({ "name": sandbox_name("s1") }))));
client
.expect_get_sandbox()
.returning(|_, id| Ok(sandbox_in_state(id, "STATE_RUNNING")));
let probes = AtomicUsize::new(0);
client.expect_execute().returning(move |_, _, _| {
let seen = probes.fetch_add(1, Ordering::SeqCst);
match answer_after {
Some(after) if seen >= after => Ok(health_reply()),
_ => Err(bad_gateway()),
}
});
client
.expect_delete_sandbox()
.times(deletes)
.returning(|_, _| Ok(()));
client
};
let ready = provider(client_with(Some(2), 0))
.create(CreateSandboxRequest::default())
.await
.expect("the agent answers on the third probe");
assert_eq!(ready.sandbox_id, "s1");
let error = provider(client_with(None, 1))
.create(CreateSandboxRequest::default())
.await
.expect_err("an agent that never answers fails create");
assert_eq!(error.code, "SANDBOX_UNREACHABLE", "{error}");
let rendered = format!("{error:?}");
assert!(
rendered.contains("Bad Gateway") && rendered.contains("did not become servable within"),
"the timeout sits over the last probe's cause: {rendered}"
);
}
#[tokio::test]
async fn a_requested_lifetime_reaches_the_create_body_and_cannot_raise_the_declared_ceiling() {
for (timeout_ms, expected) in [(90_000_u64, "90s"), (7_200_000, "3600s")] {
let mut client = MockAgentPlatformApi::new();
client
.expect_create_sandbox()
.withf(move |_, request| request.ttl.as_deref() == Some(expected))
.times(1)
.returning(|_, _| Ok(done_op(serde_json::json!({ "name": sandbox_name("s1") }))));
client
.expect_get_sandbox()
.returning(|_, id| Ok(sandbox_in_state(id, "STATE_RUNNING")));
client
.expect_execute()
.returning(|_, _, _| Ok(health_reply()));
provider(client)
.create(CreateSandboxRequest {
timeout_ms: Some(timeout_ms),
..Default::default()
})
.await
.expect("create succeeds");
}
}
#[tokio::test]
async fn create_deletes_the_sandbox_when_its_agent_never_answers() {
let mut client = MockAgentPlatformApi::new();
client
.expect_create_sandbox()
.returning(|_, _| Ok(done_op(serde_json::json!({ "name": sandbox_name("s1") }))));
client
.expect_get_sandbox()
.returning(|_, id| Ok(sandbox_in_state(id, "STATE_RUNNING")));
client
.expect_execute()
.returning(|_, _, _| Err(execute_refused()));
client
.expect_delete_sandbox()
.withf(|engine, sandbox| engine == "eng1" && sandbox == "s1")
.times(1)
.returning(|_, _| Ok(()));
provider(client)
.create(CreateSandboxRequest::default())
.await
.expect_err("a sandbox whose agent is silent is not a usable sandbox");
}
#[tokio::test]
async fn create_refuses_a_per_sandbox_environment() {
let mut client = MockAgentPlatformApi::new();
client.expect_create_sandbox().never();
let error = provider(client)
.create(CreateSandboxRequest {
env: BTreeMap::from([("TOKEN".to_string(), "secret".to_string())]),
..Default::default()
})
.await
.expect_err("a sandbox environment must be refused");
assert_eq!(error.code, "OPERATION_NOT_SUPPORTED", "{error}");
assert!(error.to_string().contains("env"), "{error}");
assert!(error.to_string().contains("per command"), "{error}");
}
#[tokio::test]
async fn get_returns_none_when_the_sandbox_is_gone() {
let mut client = MockAgentPlatformApi::new();
client
.expect_get_sandbox()
.returning(|_, _| Err(not_found()));
let found = provider(client)
.get("s1")
.await
.expect("a gone sandbox is a valid answer");
assert!(found.is_none(), "a not-found sandbox is None, not an error");
}
#[tokio::test]
async fn get_does_not_report_a_running_sandbox_whose_agent_is_silent() {
let mut client = MockAgentPlatformApi::new();
client
.expect_get_sandbox()
.returning(|_, id| Ok(sandbox_in_state(id, "STATE_RUNNING")));
client
.expect_execute()
.returning(|_, _, _| Err(execute_refused()));
let error = provider(client)
.get("s1")
.await
.expect_err("a running record with a silent agent is not a healthy sandbox");
assert_eq!(error.code, "SANDBOX_UNREACHABLE", "{error}");
}
#[tokio::test]
async fn get_or_create_replaces_a_stale_sandbox_without_deleting_it() {
let mut client = MockAgentPlatformApi::new();
client
.expect_get_sandbox()
.withf(|_, sandbox| sandbox == "stale")
.returning(|_, id| Ok(sandbox_in_state(id, "STATE_RUNNING")));
client
.expect_execute()
.withf(|_, sandbox, _| sandbox == "stale")
.returning(|_, _, _| Err(execute_refused()));
client.expect_create_sandbox().times(1).returning(|_, _| {
Ok(done_op(
serde_json::json!({ "name": sandbox_name("fresh") }),
))
});
client
.expect_get_sandbox()
.withf(|_, sandbox| sandbox == "fresh")
.returning(|_, id| Ok(sandbox_in_state(id, "STATE_RUNNING")));
client
.expect_execute()
.withf(|_, sandbox, input| sandbox == "fresh" && op_of(input) == "health")
.returning(|_, _, _| Ok(health_reply()));
client.expect_delete_sandbox().never();
let sandbox = provider(client)
.get_or_create(CreateSandboxRequest {
sandbox_id: Some("stale".to_string()),
..Default::default()
})
.await
.expect("a stale sandbox is replaced");
assert_eq!(
sandbox.sandbox.sandbox_id, "fresh",
"the fresh sandbox is returned, not the stale id"
);
assert!(sandbox.created, "a replacement is a sandbox this call made");
}
#[tokio::test]
async fn get_or_create_waits_for_a_booting_sandbox_rather_than_creating_a_second() {
let reads = Arc::new(AtomicUsize::new(0));
let mut client = MockAgentPlatformApi::new();
client.expect_get_sandbox().returning(move |_, id| {
if reads.fetch_add(1, Ordering::SeqCst) == 0 {
Ok(sandbox_in_state(id, "STATE_CREATING"))
} else {
Ok(sandbox_in_state(id, "STATE_RUNNING"))
}
});
client
.expect_execute()
.withf(|_, _, input| op_of(input) == "health")
.returning(|_, _, _| Ok(health_reply()));
client.expect_create_sandbox().never();
client.expect_delete_sandbox().never();
let sandbox = provider(client)
.get_or_create(CreateSandboxRequest {
sandbox_id: Some("booting".to_string()),
..Default::default()
})
.await
.expect("a booting sandbox is waited for and handed back");
assert_eq!(sandbox.sandbox.sandbox_id, "booting");
assert_eq!(sandbox.sandbox.state, SandboxState::Running);
assert!(
!sandbox.created,
"whoever started it created it, not this call"
);
}
#[tokio::test]
async fn get_or_create_replaces_a_paused_sandbox_without_touching_it() {
let mut client = MockAgentPlatformApi::new();
client.expect_get_sandbox().returning(|_, id| {
if id == "paused" {
Ok(sandbox_in_state(id, "STATE_PAUSED"))
} else {
Ok(sandbox_in_state(id, "STATE_RUNNING"))
}
});
client.expect_create_sandbox().times(1).returning(|_, _| {
Ok(done_op(
serde_json::json!({ "name": sandbox_name("fresh") }),
))
});
client
.expect_execute()
.withf(|_, sandbox, input| sandbox == "fresh" && op_of(input) == "health")
.returning(|_, _, _| Ok(health_reply()));
client.expect_delete_sandbox().never();
let resolved = provider(client)
.get_or_create(CreateSandboxRequest {
sandbox_id: Some("paused".to_string()),
..Default::default()
})
.await
.expect("a fresh sandbox serves the reconnect");
assert_eq!(resolved.sandbox.sandbox_id, "fresh");
assert!(resolved.created);
}
#[tokio::test]
async fn get_or_create_waits_out_a_wake_someone_else_started() {
let reads = Arc::new(AtomicUsize::new(0));
let mut client = MockAgentPlatformApi::new();
client.expect_get_sandbox().returning(move |_, id| {
if reads.fetch_add(1, Ordering::SeqCst) == 0 {
Ok(sandbox_in_state(id, "STATE_RESUMING"))
} else {
Ok(sandbox_in_state(id, "STATE_RUNNING"))
}
});
client
.expect_execute()
.withf(|_, _, input| op_of(input) == "health")
.returning(|_, _, _| Ok(health_reply()));
client.expect_create_sandbox().never();
client.expect_delete_sandbox().never();
let sandbox = provider(client)
.get_or_create(CreateSandboxRequest {
sandbox_id: Some("waking".to_string()),
..Default::default()
})
.await
.expect("a wake already in flight is waited out");
assert!(!sandbox.created, "the sandbox existed before this call");
}
async fn generation_for_boot(boot_id: &'static str) -> u64 {
let mut client = MockAgentPlatformApi::new();
client
.expect_get_sandbox()
.returning(|_, id| Ok(sandbox_in_state(id, "STATE_RUNNING")));
client
.expect_execute()
.returning(move |_, _, _| Ok(health_reply_with_boot(boot_id)));
provider(client)
.get("s1")
.await
.expect("a running sandbox")
.expect("a present sandbox")
.generation
}
#[tokio::test]
async fn generation_tracks_the_container_boot_id() {
let first = generation_for_boot("boot-id-aaaa").await;
let replaced = generation_for_boot("boot-id-bbbb").await;
let same = generation_for_boot("boot-id-aaaa").await;
assert_ne!(
first, replaced,
"a replaced container changes the generation"
);
assert_eq!(
first, same,
"the same container keeps its generation across separate reads"
);
assert_ne!(
first, NO_GENERATION,
"a probed running sandbox carries a real generation"
);
}
#[tokio::test]
async fn get_refuses_an_agent_that_reports_an_empty_boot_id() {
let mut client = MockAgentPlatformApi::new();
client
.expect_get_sandbox()
.returning(|_, id| Ok(sandbox_in_state(id, "STATE_RUNNING")));
client
.expect_execute()
.returning(|_, _, _| Ok(health_reply_with_boot("")));
let error = provider(client)
.get("s1")
.await
.expect_err("an empty boot id is no container identity");
assert_eq!(error.code, "SANDBOX_UNREACHABLE", "{error}");
}
#[tokio::test]
async fn get_refuses_an_agent_whose_health_omits_the_boot_id() {
let mut client = MockAgentPlatformApi::new();
client
.expect_get_sandbox()
.returning(|_, id| Ok(sandbox_in_state(id, "STATE_RUNNING")));
client.expect_execute().returning(|_, _, _| {
Ok(serde_json::to_vec(&serde_json::json!({ "protocolVersion": 1 })).expect("serializes"))
});
let error = provider(client)
.get("s1")
.await
.expect_err("a health reply without a boot id is not usable");
assert_eq!(error.code, "SANDBOX_UNREACHABLE", "{error}");
}
#[tokio::test(start_paused = true)]
async fn get_does_not_hang_on_a_wedged_agent() {
let error = provider_from(Arc::new(WedgedAgent::default()))
.get("s1")
.await
.expect_err("a wedged agent is unreachable, not a hang");
assert_eq!(error.code, "SANDBOX_UNREACHABLE", "{error}");
}
#[tokio::test(start_paused = true)]
async fn create_gives_up_on_a_wedged_first_probe_and_deletes_the_sandbox() {
let client = Arc::new(WedgedAgent::default());
let started = tokio::time::Instant::now();
let error = provider_from(client.clone())
.create(CreateSandboxRequest::default())
.await
.expect_err("a wedged agent fails create, not hangs it");
assert!(
started.elapsed() < AGENT_PROBE_BUDGET,
"the wait's deadline, not the probe's own budget, ended the wait"
);
assert_eq!(error.code, "SANDBOX_UNREACHABLE", "{error}");
assert!(
error.to_string().contains("did not become servable within"),
"{error}"
);
assert_eq!(client.deletes.load(Ordering::SeqCst), 1);
}
#[derive(Debug, Default)]
struct WedgedAgent {
deletes: AtomicUsize,
}
#[async_trait]
impl AgentPlatformApi for WedgedAgent {
async fn get_sandbox(&self, _engine: &str, sandbox: &str) -> ClientResult<SandboxEnvironment> {
Ok(sandbox_in_state(sandbox, "STATE_RUNNING"))
}
async fn execute(&self, _engine: &str, _sandbox: &str, _input: &[u8]) -> ClientResult<Vec<u8>> {
tokio::time::sleep(Duration::from_secs(86_400)).await;
unreachable!("the probe budget should fire before a wedged execute returns")
}
async fn create_engine(&self, _display_name: &str) -> ClientResult<Operation> {
unimplemented!()
}
async fn delete_engine(&self, _engine: &str) -> ClientResult<()> {
unimplemented!()
}
async fn list_engines(&self) -> ClientResult<Vec<ReasoningEngine>> {
unimplemented!()
}
async fn create_template(
&self,
_engine: &str,
_template: SandboxEnvironmentTemplate,
) -> ClientResult<Operation> {
unimplemented!()
}
async fn get_template(
&self,
_engine: &str,
_template: &str,
) -> ClientResult<SandboxEnvironmentTemplate> {
unimplemented!()
}
async fn delete_template(&self, _engine: &str, _template: &str) -> ClientResult<()> {
unimplemented!()
}
async fn list_templates(&self, _engine: &str) -> ClientResult<Vec<SandboxEnvironmentTemplate>> {
unimplemented!()
}
async fn create_sandbox(
&self,
_engine: &str,
_request: SandboxCreateRequest,
) -> ClientResult<Operation> {
Ok(done_op(serde_json::json!({ "name": sandbox_name("s1") })))
}
async fn list_sandboxes(&self, _engine: &str) -> ClientResult<Vec<SandboxEnvironment>> {
unimplemented!()
}
async fn delete_sandbox(&self, _engine: &str, _sandbox: &str) -> ClientResult<()> {
self.deletes.fetch_add(1, Ordering::SeqCst);
Ok(())
}
async fn snapshot(
&self,
_engine: &str,
_sandbox: &str,
_display_name: &str,
) -> ClientResult<Operation> {
unimplemented!()
}
async fn get_operation(&self, _name: &str) -> ClientResult<Operation> {
unimplemented!()
}
}
#[tokio::test]
async fn list_reports_each_sandbox_with_its_state() {
let mut client = MockAgentPlatformApi::new();
client.expect_list_sandboxes().returning(|_| {
Ok(vec![
sandbox_in_state("a", "STATE_RUNNING"),
sandbox_in_state("b", "STATE_PAUSED"),
])
});
let sandboxes = provider(client)
.list()
.await
.expect("list is supported here");
assert_eq!(sandboxes.len(), 2);
assert_eq!(sandboxes[0].sandbox_id, "a");
assert_eq!(sandboxes[0].state, SandboxState::Running);
assert_eq!(sandboxes[1].sandbox_id, "b");
assert_eq!(sandboxes[1].state, SandboxState::Paused);
}
#[tokio::test]
async fn every_state_word_the_api_sends_maps_onto_one_of_ours() {
for (reported, expected) in [
("STATE_RUNNING", SandboxState::Running),
("STATE_CREATING", SandboxState::Starting),
("STATE_PENDING", SandboxState::Starting),
("STATE_RESUMING", SandboxState::Starting),
("STATE_PAUSED", SandboxState::Paused),
("STATE_PAUSING", SandboxState::Paused),
("STATE_SUSPENDED", SandboxState::Paused),
("STATE_STOPPED", SandboxState::Terminated),
("STATE_FAILED", SandboxState::Terminated),
("STATE_DELETING", SandboxState::Terminated),
("STATE_DELETED", SandboxState::Terminated),
] {
let mut client = MockAgentPlatformApi::new();
client
.expect_list_sandboxes()
.returning(move |_| Ok(vec![sandbox_in_state("s1", reported)]));
let sandboxes = provider(client)
.list()
.await
.unwrap_or_else(|error| panic!("{reported}: {error}"));
assert_eq!(
sandboxes.len(),
1,
"{reported} was dropped as a state this provider cannot read"
);
assert_eq!(sandboxes[0].state, expected, "state {reported}");
}
}
#[tokio::test]
async fn a_state_word_outside_that_set_is_an_error_rather_than_a_guess() {
for reported in [Some("STATE_HIBERNATED"), None] {
let mut client = MockAgentPlatformApi::new();
client.expect_get_sandbox().returning(move |_, id| {
let mut sandbox = sandbox_in_state(id, "STATE_RUNNING");
sandbox.state = reported.map(str::to_string);
Ok(sandbox)
});
let error = provider(client)
.get("s1")
.await
.expect_err("an unreadable state must not become a sandbox");
assert_eq!(error.code, "UNEXPECTED_RESPONSE_FORMAT", "{error}");
}
}
#[tokio::test]
async fn the_program_leads_its_arguments_in_the_envelope() {
let mut client = MockAgentPlatformApi::new();
client
.expect_execute()
.times(1)
.withf(|_, _, input| {
let body: serde_json::Value =
serde_json::from_slice(input).expect("the envelope is json");
body["command"] == serde_json::json!(["python", "-u", "main.py"])
&& body["cwd"] == serde_json::json!("/work")
})
.returning(|_, _, _| Ok(ndjson(&[exit_frame(0)])));
let frames: Vec<_> = provider(client)
.run_command(
"s1",
RunCommandRequest {
command: "python".to_string(),
args: vec!["-u".to_string(), "main.py".to_string()],
cwd: Some("/work".to_string()),
env: BTreeMap::new(),
timeout: Duration::from_secs(5),
},
)
.await
.expect("the command runs")
.collect()
.await;
assert!(
matches!(frames.last(), Some(Ok(CommandOutput::Exit { code, .. })) if *code == 0),
"the command has to reach its exit: {frames:?}"
);
}
#[tokio::test]
async fn a_short_command_runs_synchronously_without_a_job() {
let mut client = MockAgentPlatformApi::new();
client
.expect_execute()
.returning(|_, _, input| match op_of(input).as_str() {
"exec" => Ok(ndjson(&[stdout_frame(0, b"hi"), exit_frame(0)])),
"jobStart" => panic!("a short command must not start a job"),
other => panic!("unexpected op {other}"),
});
let frames: Vec<_> = provider(client)
.run_command(
"s1",
RunCommandRequest {
command: "/bin/echo".to_string(),
args: vec!["hi".to_string()],
cwd: None,
env: BTreeMap::new(),
timeout: Duration::from_secs(5),
},
)
.await
.expect("the command runs")
.collect()
.await;
assert!(
matches!(frames.first(), Some(Ok(CommandOutput::Stdout { data, .. })) if data == b"hi")
);
assert!(matches!(
frames.last(),
Some(Ok(CommandOutput::Exit { code: 0, .. }))
));
}
#[tokio::test(start_paused = true)]
async fn a_long_command_uses_the_job_path() {
let polls = Arc::new(AtomicUsize::new(0));
let mut client = MockAgentPlatformApi::new();
client
.expect_execute()
.returning(move |_, _, input| match op_of(input).as_str() {
"jobStart" => Ok(serde_json::to_vec(&serde_json::json!({ "jobId": "j1" })).unwrap()),
"jobPoll" => {
let poll = polls.fetch_add(1, Ordering::SeqCst);
if poll == 0 {
Ok(serde_json::to_vec(&serde_json::json!({
"running": true,
"frames": [stdout_frame(0, b"work")],
}))
.unwrap())
} else {
Ok(serde_json::to_vec(&serde_json::json!({
"running": false,
"frames": [],
"exitCode": 0,
"truncated": false,
}))
.unwrap())
}
}
"exec" => panic!("a long command must not run synchronously"),
other => panic!("unexpected op {other}"),
});
let frames: Vec<_> = provider(client)
.run_command(
"s1",
RunCommandRequest {
command: "/bin/sleep".to_string(),
args: vec!["40".to_string()],
cwd: None,
env: BTreeMap::new(),
timeout: Duration::from_secs(60),
},
)
.await
.expect("the job starts")
.collect()
.await;
assert!(
matches!(frames.first(), Some(Ok(CommandOutput::Stdout { data, .. })) if data == b"work")
);
assert!(matches!(
frames.last(),
Some(Ok(CommandOutput::Exit { code: 0, .. }))
));
}
fn job_client() -> MockAgentPlatformApi {
let mut client = MockAgentPlatformApi::new();
client
.expect_execute()
.returning(|_, _, input| match op_of(input).as_str() {
"jobStart" => Ok(serde_json::to_vec(&serde_json::json!({ "jobId": "j1" })).unwrap()),
"jobPoll" => Ok(serde_json::to_vec(&match since_of(input) {
None => serde_json::json!({
"running": true,
"frames": [stdout_frame(0, b"work")],
}),
Some(_) => serde_json::json!({
"running": false,
"frames": [],
"exitCode": 0,
"truncated": false,
}),
})
.unwrap()),
other => panic!("unexpected op {other}"),
});
client
}
fn since_of(input: &[u8]) -> Option<u64> {
serde_json::from_slice::<serde_json::Value>(input)
.ok()?
.get("sinceSeq")?
.as_u64()
}
#[tokio::test(start_paused = true)]
async fn a_polled_job_carries_what_run_command_would_have_streamed() {
let streamed: Vec<CommandOutput> = provider(job_client())
.run_command("s1", long_command())
.await
.expect("the job starts")
.map(|frame| frame.expect("every frame is output"))
.collect()
.await;
let sandbox = provider(job_client());
let started = sandbox
.start_job("s1", long_command())
.await
.expect("the job starts");
assert_eq!(started.job_id, "j1");
let first = sandbox
.poll_job("s1", &started.job_id, None)
.await
.expect("the first poll answers");
assert!(first.running, "the job has not ended yet");
let last = sandbox
.poll_job("s1", &started.job_id, Some(0))
.await
.expect("the second poll answers");
assert!(!last.running);
let exit = last.exit.expect("a job that ended carries its exit");
let polled: Vec<CommandOutput> = first
.frames
.into_iter()
.chain(last.frames)
.chain([CommandOutput::Exit {
code: exit.code,
truncated: exit.truncated,
}])
.collect();
assert_eq!(polled, streamed);
assert_eq!(
polled.len(),
2,
"the fixture produces one output frame and one exit, so an empty match would prove nothing"
);
}
fn long_command() -> RunCommandRequest {
RunCommandRequest {
command: "/bin/sleep".to_string(),
args: vec!["40".to_string()],
cwd: None,
env: BTreeMap::new(),
timeout: Duration::from_secs(60),
}
}
fn cancel_watching_client(
poll: serde_json::Value,
) -> (
MockAgentPlatformApi,
tokio::sync::mpsc::UnboundedReceiver<serde_json::Value>,
) {
cancel_watching_client_answering(move || Ok(serde_json::to_vec(&poll).unwrap()))
}
fn cancel_watching_client_failing_after_one_poll(
poll: serde_json::Value,
) -> (
MockAgentPlatformApi,
tokio::sync::mpsc::UnboundedReceiver<serde_json::Value>,
) {
let polls = AtomicUsize::new(0);
cancel_watching_client_answering(move || {
if polls.fetch_add(1, Ordering::SeqCst) == 0 {
Ok(serde_json::to_vec(&poll).unwrap())
} else {
Err(execute_refused())
}
})
}
fn cancel_refusing_client(
poll: serde_json::Value,
) -> (
MockAgentPlatformApi,
tokio::sync::mpsc::UnboundedReceiver<serde_json::Value>,
) {
let (cancels, seen) = tokio::sync::mpsc::unbounded_channel();
let mut client = MockAgentPlatformApi::new();
client
.expect_execute()
.returning(move |_, _, input| match op_of(input).as_str() {
"jobStart" => Ok(serde_json::to_vec(&serde_json::json!({ "jobId": "j1" })).unwrap()),
"jobPoll" => Ok(serde_json::to_vec(&poll).unwrap()),
"jobCancel" => {
cancels
.send(serde_json::from_slice(input).expect("a cancel body is json"))
.expect("the test still listens for cancels");
Err(execute_refused())
}
other => panic!("unexpected op {other}"),
});
(client, seen)
}
fn cancel_watching_client_answering(
mut poll: impl FnMut() -> ClientResult<Vec<u8>> + Send + 'static,
) -> (
MockAgentPlatformApi,
tokio::sync::mpsc::UnboundedReceiver<serde_json::Value>,
) {
let (cancels, seen) = tokio::sync::mpsc::unbounded_channel();
let mut client = MockAgentPlatformApi::new();
client
.expect_execute()
.returning(move |_, _, input| match op_of(input).as_str() {
"jobStart" => Ok(serde_json::to_vec(&serde_json::json!({ "jobId": "j1" })).unwrap()),
"jobPoll" => poll(),
"jobCancel" => {
cancels
.send(serde_json::from_slice(input).expect("a cancel body is json"))
.expect("the test still listens for cancels");
Ok(b"{}".to_vec())
}
other => panic!("unexpected op {other}"),
});
(client, seen)
}
#[tokio::test(start_paused = true)]
async fn dropping_a_running_jobs_stream_cancels_the_job() {
let (client, mut cancels) = cancel_watching_client(serde_json::json!({
"running": true,
"frames": [stdout_frame(0, b"work")],
}));
let sandbox = provider(client);
let mut frames = sandbox
.run_command("s1", long_command())
.await
.expect("the job starts");
assert!(
matches!(frames.next().await, Some(Ok(CommandOutput::Stdout { .. }))),
"the job is running and has produced output"
);
drop(frames);
let cancel = tokio::time::timeout(Duration::from_secs(5), cancels.recv())
.await
.expect("the drop cancels the job")
.expect("the cancel carries a body");
assert_eq!(cancel["op"], "jobCancel", "{cancel}");
assert_eq!(cancel["jobId"], "j1", "{cancel}");
}
#[tokio::test(start_paused = true)]
async fn a_job_that_reported_its_exit_is_not_cancelled_afterwards() {
let (client, mut cancels) = cancel_watching_client(serde_json::json!({
"running": false,
"frames": [stdout_frame(0, b"work")],
"exitCode": 0,
"truncated": false,
}));
let sandbox = provider(client);
let frames: Vec<_> = sandbox
.run_command("s1", long_command())
.await
.expect("the job starts")
.collect()
.await;
assert!(
matches!(frames.last(), Some(Ok(CommandOutput::Exit { code: 0, .. }))),
"the job ended on its own"
);
assert!(
tokio::time::timeout(Duration::from_secs(5), cancels.recv())
.await
.is_err(),
"a job that ended on its own is not cancelled afterwards"
);
}
#[tokio::test(start_paused = true)]
async fn a_stream_that_ended_on_a_failed_poll_still_cancels_the_job() {
let (client, mut cancels) = cancel_watching_client_failing_after_one_poll(serde_json::json!({
"running": true,
"frames": [stdout_frame(0, b"work")],
}));
let sandbox = provider(client);
let frames: Vec<_> = sandbox
.run_command("s1", long_command())
.await
.expect("the job starts")
.collect()
.await;
let Some(Err(error)) = frames.last() else {
panic!("the failed poll reaches the caller: {frames:?}");
};
assert!(
error.to_string().contains("the job is no longer watched"),
"{error}"
);
let cancel = tokio::time::timeout(Duration::from_secs(5), cancels.recv())
.await
.expect("the drop cancels the job the failed poll left running")
.expect("the cancel carries a body");
assert_eq!(cancel["op"], "jobCancel", "{cancel}");
assert_eq!(cancel["jobId"], "j1", "{cancel}");
}
#[tokio::test(start_paused = true)]
async fn a_deadline_arms_the_drop_unless_its_cancel_landed() {
let (client, mut cancels) = cancel_refusing_client(serde_json::json!({
"running": true,
"frames": [],
}));
let refused = provider(client);
let frames: Vec<_> = refused
.run_command("s1", long_command())
.await
.expect("the job starts")
.collect()
.await;
let Some(Err(error)) = frames.last() else {
panic!("the elapsed deadline reaches the caller: {frames:?}");
};
assert_eq!(error.code, "SANDBOX_OUTCOME_UNKNOWN", "{error}");
assert!(
error.to_string().contains("could not be cancelled"),
"{error}"
);
for attempt in ["the deadline", "the drop"] {
let cancel = tokio::time::timeout(Duration::from_secs(5), cancels.recv())
.await
.unwrap_or_else(|_| panic!("{attempt} cancels the job it could not stop"))
.expect("the cancel carries a body");
assert_eq!(cancel["op"], "jobCancel", "{cancel}");
assert_eq!(cancel["jobId"], "j1", "{cancel}");
}
let (client, mut cancels) = cancel_watching_client(serde_json::json!({
"running": true,
"frames": [],
}));
let landed = provider(client);
let frames: Vec<_> = landed
.run_command("s1", long_command())
.await
.expect("the job starts")
.collect()
.await;
let Some(Err(error)) = frames.last() else {
panic!("the elapsed deadline reaches the caller: {frames:?}");
};
assert!(error.to_string().contains("timeoutExceeded"), "{error}");
tokio::time::timeout(Duration::from_secs(5), cancels.recv())
.await
.expect("the deadline cancels the job")
.expect("the cancel carries a body");
assert!(
tokio::time::timeout(Duration::from_secs(5), cancels.recv())
.await
.is_err(),
"a job whose cancel landed is not cancelled again on drop"
);
}
#[tokio::test(start_paused = true)]
async fn a_job_error_object_becomes_a_stream_error() {
let mut client = MockAgentPlatformApi::new();
client
.expect_execute()
.returning(|_, _, input| match op_of(input).as_str() {
"jobStart" => Ok(serde_json::to_vec(&serde_json::json!({ "jobId": "j1" })).unwrap()),
"jobPoll" => Ok(serde_json::to_vec(&serde_json::json!({
"running": false,
"frames": [],
"error": { "code": "timeoutExceeded", "message": "exceeded its 60000ms timeout" },
}))
.unwrap()),
other => panic!("unexpected op {other}"),
});
let frames: Vec<_> = provider(client)
.run_command(
"s1",
RunCommandRequest {
command: "/bin/sleep".to_string(),
args: vec!["99".to_string()],
cwd: None,
env: BTreeMap::new(),
timeout: Duration::from_secs(60),
},
)
.await
.expect("the job starts")
.collect()
.await;
let error = frames
.last()
.expect("a terminal item")
.as_ref()
.expect_err("an error object is a failure");
assert!(error.to_string().contains("timeoutExceeded"), "{error}");
}
#[tokio::test(start_paused = true)]
async fn a_failed_command_is_delivered_once_where_a_read_still_retries() {
let reads = Arc::new(AtomicUsize::new(0));
let reads_seen = reads.clone();
let mut client = MockAgentPlatformApi::new();
client
.expect_execute()
.times(1)
.returning(|_, _, _| Err(execute_refused()));
client
.expect_delete_sandbox()
.times(1)
.returning(|_, _| Ok(()));
client.expect_get_sandbox().returning(move |_, id| {
if reads.fetch_add(1, Ordering::SeqCst) < 2 {
Ok(sandbox_in_state(id, "STATE_RUNNING"))
} else {
Err(not_found())
}
});
let sut = provider(client);
let Err(command) = sut
.run_command(
"s1",
RunCommandRequest {
command: "/bin/true".to_string(),
args: Vec::new(),
cwd: None,
env: BTreeMap::new(),
timeout: Duration::from_secs(5),
},
)
.await
else {
panic!("a mutating command whose execute fails is refused, not retried into success");
};
assert_eq!(command.code, "SANDBOX_OUTCOME_UNKNOWN", "{command}");
assert!(
!command.retryable,
"the call may have started the command, so it is delivered once: {command}"
);
sut.terminate("s1")
.await
.expect("the poll confirms the sandbox is gone");
assert!(
reads_seen.load(Ordering::SeqCst) > 1,
"confirming the sandbox gone took more than one read, so the read path retries"
);
}
#[tokio::test]
async fn a_command_on_a_gone_sandbox_is_refused_and_deletes_nothing() {
let mut client = MockAgentPlatformApi::new();
client
.expect_execute()
.returning(|_, _, _| Err(not_found()));
client.expect_delete_sandbox().never();
let Err(error) = provider(client)
.run_command(
"s1",
RunCommandRequest {
command: "/bin/true".to_string(),
args: Vec::new(),
cwd: None,
env: BTreeMap::new(),
timeout: Duration::from_secs(5),
},
)
.await
else {
panic!("a command against a gone sandbox is refused");
};
assert_eq!(error.code, "SANDBOX_COMMAND_FAILED", "{error}");
assert!(error.to_string().contains("sandboxGone"), "{error}");
}
#[tokio::test]
async fn a_command_without_a_timeout_or_program_is_refused() {
let Err(empty) = provider(MockAgentPlatformApi::new())
.run_command(
"s1",
RunCommandRequest {
command: String::new(),
args: Vec::new(),
cwd: None,
env: BTreeMap::new(),
timeout: Duration::from_secs(5),
},
)
.await
else {
panic!("an empty command is refused");
};
assert_eq!(empty.code, "INVALID_INPUT", "{empty}");
let Err(zero) = provider(MockAgentPlatformApi::new())
.run_command(
"s1",
RunCommandRequest {
command: "/bin/true".to_string(),
args: Vec::new(),
cwd: None,
env: BTreeMap::new(),
timeout: Duration::ZERO,
},
)
.await
else {
panic!("a zero timeout is refused");
};
assert!(zero.to_string().contains("timeout"), "{zero}");
}
#[tokio::test]
async fn write_files_sends_contents_base64_and_accepts_an_empty_body() {
let mut client = MockAgentPlatformApi::new();
client
.expect_execute()
.withf(|_, _, input| {
let value: serde_json::Value = serde_json::from_slice(input).unwrap();
op_of(input) == "writeFile"
&& value.get("contentsBase64").and_then(|v| v.as_str())
== Some(&BASE64.encode(b"data"))
&& value.get("contents").is_none()
})
.times(1)
.returning(|_, _, _| Ok(Vec::new()));
provider(client)
.write_files(
"s1",
BTreeMap::from([("a.txt".to_string(), b"data".to_vec())]),
)
.await
.expect("an empty body is a successful write");
}
#[tokio::test]
async fn read_file_decodes_the_agent_reply() {
let mut client = MockAgentPlatformApi::new();
client
.expect_execute()
.withf(|_, _, input| op_of(input) == "readFile")
.returning(|_, _, _| {
Ok(serde_json::to_vec(
&serde_json::json!({ "contentsBase64": BASE64.encode(b"file body") }),
)
.unwrap())
});
let contents = provider(client)
.read_file("s1", "a.txt")
.await
.expect("read succeeds");
assert_eq!(contents, b"file body");
}
#[tokio::test]
async fn a_relayed_agent_404_is_a_refusal_not_a_gone_sandbox() {
fn relayed(status: u16, details: &'static str) -> AlienError<AgentPlatformErrorData> {
let body = serde_json::json!({
"error": {
"code": status,
"message": format!("Execution Failed. URL not found `https://x.sandbox.vertexai.goog`. Error Details: {details}"),
}
})
.to_string();
let http = AlienError::new(alien_client_core::ErrorData::HttpResponseError {
message: format!("Request failed with HTTP {status}"),
url: "https://example.invalid/:execute".to_string(),
http_status: status,
http_request_text: None,
http_response_text: Some(body),
});
let http = if status == 404 {
http.context(alien_client_core::ErrorData::RemoteResourceNotFound {
resource_type: "Vertex AI Agent Platform".to_string(),
resource_name: "s1".to_string(),
})
} else {
http
};
http.context(AgentPlatformErrorData::ExecuteFailed {
sandbox: "s1".to_string(),
message: "the API rejected or cut short the request".to_string(),
})
}
fn relayed_404(details: &'static str) -> AlienError<AgentPlatformErrorData> {
relayed(404, details)
}
let mut client = MockAgentPlatformApi::new();
client.expect_execute().returning(|_, _, _| {
Err(relayed_404(
"PATH_NOT_FOUND: No such file in the sandbox: /before",
))
});
let error = provider(client)
.read_file("s1", "/before")
.await
.expect_err("a missing file is an error");
let rendered = format!("{error:?}");
assert!(rendered.contains("agentRefused"), "{rendered}");
assert!(rendered.contains("PATH_NOT_FOUND"), "{rendered}");
assert!(agent_answer(¬_found()).is_none());
assert!(agent_answer(&relayed_404(
"Bad Gateway: Unable to reach the sandbox environment."
))
.is_none());
let mut client = MockAgentPlatformApi::new();
client
.expect_execute()
.returning(move |_, _, _| Err(relayed(500, "INTERNAL_SERVER_ERROR: stream reset")));
let Err(command) = provider(client)
.run_command(
"s1",
RunCommandRequest {
command: "/bin/true".to_string(),
args: Vec::new(),
cwd: None,
env: BTreeMap::new(),
timeout: Duration::from_secs(5),
},
)
.await
else {
panic!("a relayed 5xx fails the command");
};
assert_eq!(command.code, "SANDBOX_OUTCOME_UNKNOWN", "{command}");
}
#[tokio::test]
async fn pause_and_resume_are_refused_as_unsupported() {
assert!(!SandboxCapabilities::gcp_agent_platform().pause_resume);
let provider = provider(MockAgentPlatformApi::new());
for error in [
provider.pause("s1").await.expect_err("pause is refused"),
provider.resume("s1").await.expect_err("resume is refused"),
] {
assert_eq!(error.code, "OPERATION_NOT_SUPPORTED", "{error}");
}
}
#[tokio::test]
async fn snapshot_returns_the_snapshot_name() {
let mut client = MockAgentPlatformApi::new();
let name =
"projects/p/locations/us-central1/reasoningEngines/eng1/sandboxEnvironmentSnapshots/snap1";
client
.expect_snapshot()
.withf(|engine, sandbox, display| {
engine == "eng1" && sandbox == "s1" && !display.is_empty()
})
.returning(move |_, _, _| Ok(done_op(serde_json::json!({ "name": name }))));
let returned = provider(client)
.snapshot("s1")
.await
.expect("snapshot completes");
assert_eq!(returned, name);
}
#[tokio::test(start_paused = true)]
async fn terminate_confirms_by_polling_to_not_found() {
let reads = Arc::new(AtomicUsize::new(0));
let mut client = MockAgentPlatformApi::new();
client
.expect_delete_sandbox()
.times(1)
.returning(|_, _| Ok(()));
client.expect_get_sandbox().returning(move |_, id| {
if reads.fetch_add(1, Ordering::SeqCst) == 0 {
Ok(sandbox_in_state(id, "STATE_RUNNING"))
} else {
Err(not_found())
}
});
provider(client)
.terminate("s1")
.await
.expect("a sandbox that goes absent is confirmed gone");
}
#[tokio::test(start_paused = true)]
async fn terminate_reports_unconfirmed_when_the_sandbox_stays_present() {
let mut client = MockAgentPlatformApi::new();
client.expect_delete_sandbox().returning(|_, _| Ok(()));
client
.expect_get_sandbox()
.returning(|_, id| Ok(sandbox_in_state(id, "STATE_RUNNING")));
let error = provider(client)
.terminate("s1")
.await
.expect_err("a sandbox still present after the poll is not contained");
assert!(
error.to_string().contains("may still be running"),
"{error}"
);
}
#[test]
fn egress_refuses_domain_scoping_and_names_the_modes() {
let error = egress_control_config(
"sbx-7",
&SandboxEgress::AllowDomains {
domains: vec!["x.io".into()],
},
)
.expect_err("domain-scoped egress has no representation");
assert_eq!(error.code, "INVALID_INPUT", "{error}");
let rendered = error.to_string();
assert!(rendered.contains("sbx-7"), "names the sandbox: {rendered}");
assert!(
rendered.contains("allow") && rendered.contains("deny"),
"names both modes: {rendered}"
);
assert_eq!(
egress_control_config("s", &SandboxEgress::Deny)
.expect("deny maps")
.internet_access,
Some(false)
);
assert_eq!(
egress_control_config("s", &SandboxEgress::Allow)
.expect("allow maps")
.internet_access,
Some(true)
);
}
#[tokio::test]
async fn a_sandbox_id_that_could_escape_its_sandbox_is_refused() {
for id in [
"../other",
"a/b",
"has space",
"",
"with?query",
"with#frag",
] {
let error = provider(MockAgentPlatformApi::new())
.get(id)
.await
.expect_err(&format!("'{id}' must be refused before it reaches a URL"));
assert_eq!(error.code, "INVALID_INPUT", "'{id}': {error}");
}
}
#[test]
fn an_output_without_a_terminal_frame_is_an_unknown_outcome() {
let frames = parse_exec_frames(&ndjson(&[stdout_frame(0, b"partial")])).expect("frames parse");
assert_eq!(frames.len(), 2);
frames[0].as_ref().expect("the stdout frame still arrives");
let error = frames[1]
.as_ref()
.expect_err("a truncated stream is not success");
assert!(
error.to_string().contains("without a terminal frame"),
"{error}"
);
assert_eq!(error.code, "SANDBOX_OUTCOME_UNKNOWN", "{error}");
assert!(
!error.retryable,
"the command started and its end was lost, so a repeat would run it twice: {error}"
);
}
#[test]
fn a_frame_that_does_not_decode_leaves_the_outcome_unknown() {
let frames = parse_exec_frames(&ndjson(&[serde_json::json!({
"t": "stdout",
"seq": 0,
"data": "!!not base64!!",
})]))
.expect("frames parse");
let error = frames[0]
.as_ref()
.expect_err("a payload that does not decode is not output");
assert_eq!(error.code, "SANDBOX_OUTCOME_UNKNOWN", "{error}");
assert!(
error.to_string().contains("base64"),
"the decode failure must stay in the chain: {error}"
);
}
#[test]
fn a_frame_that_does_not_parse_after_output_leaves_the_outcome_unknown() {
let mut body = ndjson(&[stdout_frame(0, b"partial")]);
body.extend_from_slice(b"{not json at all}\n");
let frames = parse_exec_frames(&body).expect("frames parse");
let error = frames
.last()
.expect("a trailing item")
.as_ref()
.expect_err("a malformed frame is not output");
assert_eq!(error.code, "SANDBOX_OUTCOME_UNKNOWN", "{error}");
}
#[test]
fn a_frame_that_does_not_convert_ends_the_body() {
let mut body = ndjson(&[serde_json::json!({
"t": "stdout",
"seq": 0,
"data": "!!not base64!!",
})]);
body.extend_from_slice(&ndjson(&[exit_frame(0)]));
let frames = parse_exec_frames(&body).expect("frames parse");
assert_eq!(
frames.len(),
1,
"the exit frame must not follow the failure"
);
assert_eq!(
frames[0]
.as_ref()
.expect_err("a bad payload is not output")
.code,
"SANDBOX_OUTCOME_UNKNOWN"
);
}
#[test]
fn a_non_frame_body_is_reported_as_a_refusal() {
let error = parse_exec_frames(b"forbidden: a capability is required")
.expect_err("an error body is not a stream");
assert_eq!(error.code, "SANDBOX_COMMAND_FAILED", "{error}");
}
#[test]
fn not_found_is_read_from_the_source_chain() {
assert!(
is_not_found(¬_found()),
"a wrapped 404 is a gone sandbox"
);
assert!(
!is_not_found(&execute_refused()),
"an ordinary execute failure is not a gone sandbox"
);
}
#[test]
fn the_engine_is_reduced_to_a_bare_segment() {
let provider = provider(MockAgentPlatformApi::new());
assert_eq!(
provider.engine(),
"eng1",
"the full resource name is reduced to the engine id"
);
}
fn access_denied() -> AlienError<AgentPlatformErrorData> {
AlienError::new(alien_client_core::ErrorData::RemoteAccessDenied {
resource_type: "SandboxEnvironment".to_string(),
resource_name: "s1".to_string(),
})
.context(AgentPlatformErrorData::RequestFailed {
operation: "delete sandbox".to_string(),
message: "s1".to_string(),
})
}
#[test]
fn capabilities_describe_the_backend_not_this_declaration() {
let platform = SandboxCapabilities::gcp_agent_platform();
assert_eq!(
provider(MockAgentPlatformApi::new()).capabilities(),
platform
);
let untimed = GcpAgentPlatformSandbox::new(
Arc::new(MockAgentPlatformApi::new()),
ENGINE_FULL.to_string(),
TEMPLATE.to_string(),
None,
);
assert_eq!(
untimed.capabilities(),
platform,
"a sandbox with no declared ttl still expires, so the row does not change"
);
assert!(
platform.sandbox_lifetime,
"Agent Platform always sets expireTime"
);
}
#[tokio::test]
async fn a_tenant_key_is_refused_rather_than_dropped() {
let mut client = MockAgentPlatformApi::new();
client.expect_create_sandbox().never();
let error = provider(client)
.create(CreateSandboxRequest {
tenant_key: Some("tenant-1".to_string()),
..Default::default()
})
.await
.expect_err("a tenant key Agent Platform cannot honour is refused");
assert_eq!(error.code, "OPERATION_NOT_SUPPORTED", "{error}");
assert!(
error.to_string().contains("tenantKey"),
"the refusal has to name the field a caller must remove: {error}"
);
}
#[tokio::test]
async fn terminate_of_an_absent_sandbox_succeeds() {
let mut client = MockAgentPlatformApi::new();
client
.expect_delete_sandbox()
.times(1)
.returning(|_, _| Err(not_found()));
client.expect_get_sandbox().never();
provider(client)
.terminate("s1")
.await
.expect("terminating an absent sandbox succeeds");
}
#[tokio::test]
async fn terminate_of_a_sandbox_this_deployment_cannot_reach_is_still_refused() {
let mut client = MockAgentPlatformApi::new();
client
.expect_delete_sandbox()
.times(1)
.returning(|_, _| Err(access_denied()));
client.expect_get_sandbox().never();
let error = provider(client)
.terminate("s1")
.await
.expect_err("a refused delete is not containment");
assert_eq!(error.code, "SANDBOX_UNREACHABLE", "{error}");
assert!(
format!("{error}").to_lowercase().contains("denied"),
"the refusal must carry why the delete was refused, got: {error}"
);
}