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 session = provider(client)
.create(CreateSessionRequest::default())
.await
.expect("create succeeds");
assert_eq!(session.session_id, "s1");
assert_eq!(session.state, SandboxSessionState::Running);
}
#[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(CreateSessionRequest::default())
.await
.expect_err("a sandbox whose agent is silent is not a usable session");
}
#[tokio::test]
async fn create_refuses_a_per_session_environment() {
let mut client = MockAgentPlatformApi::new();
client.expect_create_sandbox().never();
let error = provider(client)
.create(CreateSessionRequest {
env: BTreeMap::from([("TOKEN".to_string(), "secret".to_string())]),
..Default::default()
})
.await
.expect_err("a session environment must be refused");
assert_eq!(error.code, "INVALID_INPUT", "{error}");
assert!(error.to_string().contains("each 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_session_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 session");
assert_eq!(error.code, "SANDBOX_UNREACHABLE", "{error}");
}
#[tokio::test]
async fn get_or_create_replaces_a_stale_session_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 session = provider(client)
.get_or_create(CreateSessionRequest {
session_id: Some("stale".to_string()),
..Default::default()
})
.await
.expect("a stale session is replaced");
assert_eq!(
session.session_id, "fresh",
"the fresh session is returned, not the stale id"
);
}
#[tokio::test]
async fn get_or_create_resumes_a_suspended_session_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_PAUSED"))
} else {
Ok(sandbox_in_state(id, "STATE_RUNNING"))
}
});
client
.expect_resume()
.times(1)
.returning(|_, _| Ok(done_op(serde_json::json!({}))));
client
.expect_execute()
.withf(|_, _, input| op_of(input) == "health")
.returning(|_, _, _| Ok(health_reply()));
client.expect_create_sandbox().never();
client.expect_delete_sandbox().never();
let session = provider(client)
.get_or_create(CreateSessionRequest {
session_id: Some("paused".to_string()),
..Default::default()
})
.await
.expect("a suspended session is resumed and returned");
assert_eq!(session.session_id, "paused");
assert_eq!(session.state, SandboxSessionState::Running);
assert_ne!(
session.generation, NO_GENERATION,
"a woken session carries its container generation"
);
}
#[tokio::test]
async fn get_or_create_fails_rather_than_leaking_a_resume_it_cannot_roll_back() {
let mut client = MockAgentPlatformApi::new();
client
.expect_get_sandbox()
.returning(|_, id| Ok(sandbox_in_state(id, "STATE_PAUSED")));
client
.expect_resume()
.times(1)
.returning(|_, _| Ok(done_op(serde_json::json!({}))));
client
.expect_pause()
.times(1)
.returning(|_, _| Err(not_found()));
client
.expect_execute()
.returning(|_, _, _| Ok(health_reply()));
client.expect_create_sandbox().never();
client.expect_delete_sandbox().never();
let error = provider(client)
.get_or_create(CreateSessionRequest {
session_id: Some("paused".to_string()),
..Default::default()
})
.await
.expect_err("a resume that cannot be rolled back must fail, not leak a live session");
assert!(
error.to_string().contains("paused"),
"the failure names the woken session so it stays identifiable: {error}"
);
}
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 session")
.expect("a present session")
.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 session 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))
.get("s1")
.await
.expect_err("a wedged agent is unreachable, not a hang");
assert_eq!(error.code, "SANDBOX_UNREACHABLE", "{error}");
}
#[derive(Debug)]
struct WedgedAgent;
#[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> {
unimplemented!()
}
async fn list_sandboxes(&self, _engine: &str) -> ClientResult<Vec<SandboxEnvironment>> {
unimplemented!()
}
async fn delete_sandbox(&self, _engine: &str, _sandbox: &str) -> ClientResult<()> {
unimplemented!()
}
async fn pause(&self, _engine: &str, _sandbox: &str) -> ClientResult<Operation> {
unimplemented!()
}
async fn resume(&self, _engine: &str, _sandbox: &str) -> ClientResult<Operation> {
unimplemented!()
}
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_maps_sandboxes_to_sessions() {
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 sessions = provider(client)
.list()
.await
.expect("list is supported here");
assert_eq!(sessions.len(), 2);
assert_eq!(sessions[0].session_id, "a");
assert_eq!(sessions[0].state, SandboxSessionState::Running);
assert_eq!(sessions[1].session_id, "b");
assert_eq!(sessions[1].state, SandboxSessionState::Suspended);
}
#[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: vec!["/bin/echo".to_string(), "hi".to_string()],
working_directory: None,
env: BTreeMap::new(),
deadline: 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: vec!["/bin/sleep".to_string(), "40".to_string()],
working_directory: None,
env: BTreeMap::new(),
deadline: 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, .. }))
));
}
#[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": "deadlineExceeded", "message": "exceeded its 60000ms deadline" },
}))
.unwrap()),
other => panic!("unexpected op {other}"),
});
let frames: Vec<_> = provider(client)
.run_command(
"s1",
RunCommandRequest {
command: vec!["/bin/sleep".to_string(), "99".to_string()],
working_directory: None,
env: BTreeMap::new(),
deadline: 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("deadlineExceeded"), "{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: vec!["/bin/true".to_string()],
working_directory: None,
env: BTreeMap::new(),
deadline: Duration::from_secs(5),
},
)
.await
else {
panic!("a mutating command whose execute fails is refused, not retried into success");
};
assert_eq!(command.code, "SANDBOX_COMMAND_FAILED", "{command}");
sut.terminate("s1")
.await
.expect("the poll confirms the session is gone");
assert!(
reads_seen.load(Ordering::SeqCst) > 1,
"confirming the session gone took more than one read, so the read path retries"
);
}
#[tokio::test]
async fn a_command_on_a_gone_session_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: vec!["/bin/true".to_string()],
working_directory: None,
env: BTreeMap::new(),
deadline: Duration::from_secs(5),
},
)
.await
else {
panic!("a command against a gone session is refused");
};
assert_eq!(error.code, "SANDBOX_COMMAND_FAILED", "{error}");
assert!(error.to_string().contains("sessionGone"), "{error}");
}
#[tokio::test]
async fn a_command_without_a_deadline_or_program_is_refused() {
let Err(empty) = provider(MockAgentPlatformApi::new())
.run_command(
"s1",
RunCommandRequest {
command: vec![],
working_directory: None,
env: BTreeMap::new(),
deadline: 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: vec!["/bin/true".to_string()],
working_directory: None,
env: BTreeMap::new(),
deadline: Duration::ZERO,
},
)
.await
else {
panic!("a zero deadline is refused");
};
assert!(zero.to_string().contains("deadline"), "{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 mkdir_accepts_an_empty_body() {
let mut client = MockAgentPlatformApi::new();
client
.expect_execute()
.withf(|_, _, input| op_of(input) == "mkdir")
.returning(|_, _, _| Ok(Vec::new()));
provider(client)
.mkdir("s1", "out")
.await
.expect("mkdir succeeds on an empty body");
}
#[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 suspend_and_resume_await_their_operations() {
let mut client = MockAgentPlatformApi::new();
client
.expect_pause()
.times(1)
.returning(|_, _| Ok(done_op(serde_json::json!({}))));
client
.expect_resume()
.times(1)
.returning(|_, _| Ok(done_op(serde_json::json!({}))));
let provider = provider(client);
provider.suspend("s1").await.expect("suspend completes");
provider.resume("s1").await.expect("resume completes");
}
#[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 session that goes absent is confirmed gone");
}
#[tokio::test(start_paused = true)]
async fn terminate_reports_unconfirmed_when_the_session_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 session 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_session_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}"
);
}
#[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 session"
);
assert!(
!is_not_found(&execute_refused()),
"an ordinary execute failure is not a gone session"
);
}
#[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"
);
}