1use std::collections::{BTreeMap, VecDeque};
10use std::sync::Arc;
11use std::time::Duration;
12
13use async_trait::async_trait;
14use base64::engine::general_purpose::STANDARD as BASE64;
15use base64::Engine as _;
16use futures::stream::{self, BoxStream};
17use serde::Deserialize;
18use serde_json::json;
19use tracing::warn;
20
21use crate::error::{ErrorData, Result};
22use crate::traits::{
23 Binding, CommandOutput, CreateSandboxRequest, JobError, JobExit, JobPoll, JobStart,
24 PreviewCapability, ResolvedSandbox, RunCommandRequest, Sandbox, SandboxInstance, SandboxState,
25};
26use alien_core::{SandboxCapabilities, SandboxEgress};
27use alien_error::{AlienError, Context, ContextError};
28use alien_gcp_clients::gcp::agent_platform::{
29 AgentPlatformApi, AgentPlatformErrorData, EgressControlConfig, SandboxCreateRequest,
30 SandboxEnvironment, SandboxSnapshot,
31};
32use alien_gcp_clients::gcp::longrunning::{Operation, OperationResult};
33
34const AGENT_PROTOCOL_VERSION: u32 = 1;
37
38const MAX_SYNCHRONOUS_TIMEOUT: Duration = Duration::from_secs(30);
43
44const MAX_SANDBOX_ID: usize = 63;
47
48const SANDBOX_READY_ATTEMPTS: u32 = 150;
50const SANDBOX_READY_INTERVAL: Duration = Duration::from_secs(2);
51
52#[cfg(not(test))]
55const AGENT_READY_TIMEOUT: Duration = Duration::from_secs(60);
56#[cfg(not(test))]
57const AGENT_READY_POLL: Duration = Duration::from_secs(1);
58#[cfg(test)]
60const AGENT_READY_TIMEOUT: Duration = Duration::from_millis(50);
61#[cfg(test)]
62const AGENT_READY_POLL: Duration = Duration::from_millis(5);
63
64const OPERATION_POLL_ATTEMPTS: u32 = 150;
67const OPERATION_POLL_INTERVAL: Duration = Duration::from_secs(2);
68
69const TERMINATE_POLL_ATTEMPTS: u32 = 30;
72const TERMINATE_POLL_INTERVAL: Duration = Duration::from_secs(2);
73
74const JOB_POLL_INTERVAL: Duration = Duration::from_secs(1);
76
77const JOB_POLL_GRACE: Duration = Duration::from_secs(15);
81
82const CREATE: &str = "sandbox.create";
83const GET: &str = "sandbox.get";
84const GET_OR_CREATE: &str = "sandbox.getOrCreate";
85const RUN_COMMAND: &str = "sandbox.runCommand";
86const JOB_START: &str = "sandbox.jobStart";
87const JOB_POLL: &str = "sandbox.jobPoll";
88const JOB_CANCEL: &str = "sandbox.jobCancel";
89const TERMINATE: &str = "sandbox.terminate";
90
91const NO_GENERATION: u64 = 0;
96
97const AGENT_PROBE_BUDGET: Duration = Duration::from_secs(60);
104
105pub fn egress_control_config(
114 sandbox_label: &str,
115 egress: &SandboxEgress,
116) -> Result<EgressControlConfig> {
117 let Some(internet_access) = egress.internet_access_switch() else {
118 return Err(AlienError::new(ErrorData::InvalidInput {
119 operation_context: "sandbox.template".to_string(),
120 details: format!(
121 "sandbox '{sandbox_label}' asked for domain-scoped egress, which Agent \
122 Platform cannot express; it offers only 'allow' (open) and 'deny' (closed)"
123 ),
124 field_name: Some("egress".to_string()),
125 }));
126 };
127
128 Ok(EgressControlConfig {
129 internet_access: Some(internet_access),
130 extra: Default::default(),
131 })
132}
133
134#[derive(Debug)]
136pub struct GcpAgentPlatformSandbox {
137 client: Arc<dyn AgentPlatformApi>,
138 engine: String,
141 template: String,
143 max_lifetime_seconds: Option<u32>,
145}
146
147impl GcpAgentPlatformSandbox {
148 fn lifetime_seconds(&self, timeout_ms: Option<u64>, operation: &str) -> Result<Option<u32>> {
151 match timeout_ms {
152 Some(timeout_ms) => {
153 super::requested_lifetime_seconds(timeout_ms, self.max_lifetime_seconds, operation)
154 .map(Some)
155 }
156 None => Ok(self.max_lifetime_seconds),
157 }
158 }
159
160 pub fn new(
165 client: Arc<dyn AgentPlatformApi>,
166 engine: String,
167 template: String,
168 max_lifetime_seconds: Option<u32>,
169 ) -> Self {
170 let engine = engine.rsplit('/').next().unwrap_or(&engine).to_string();
171 Self {
172 client,
173 engine,
174 template,
175 max_lifetime_seconds,
176 }
177 }
178
179 #[cfg(test)]
182 pub(crate) fn engine(&self) -> &str {
183 &self.engine
184 }
185
186 fn unsupported(&self, capability: &str, reason: &str) -> AlienError<ErrorData> {
187 AlienError::new(ErrorData::OperationNotSupported {
188 operation: capability.to_string(),
189 reason: reason.to_string(),
190 })
191 }
192
193 fn pause_resume_unsupported(&self) -> AlienError<ErrorData> {
196 self.unsupported(
197 alien_core::SandboxCapability::PauseResume.as_str(),
198 "Agent Platform sandboxes cannot be paused and resumed with their state kept",
199 )
200 }
201
202 fn checked_sandbox_id(operation: &str, sandbox_id: &str) -> Result<()> {
208 if is_addressable_id(sandbox_id) {
209 return Ok(());
210 }
211 Err(AlienError::new(ErrorData::InvalidInput {
212 operation_context: operation.to_string(),
213 details: format!(
214 "sandbox id '{sandbox_id}' must be a single segment of letters, digits, '-' and \
215 '_', at most {MAX_SANDBOX_ID} characters"
216 ),
217 field_name: Some("sandboxId".to_string()),
218 }))
219 }
220
221 async fn read_sandbox(
223 &self,
224 operation: &str,
225 sandbox_id: &str,
226 ) -> Result<Option<SandboxEnvironment>> {
227 match self.client.get_sandbox(&self.engine, sandbox_id).await {
228 Ok(sandbox) => Ok(Some(sandbox)),
229 Err(error) if is_not_found(&error) => Ok(None),
230 Err(error) => Err(error.context(ErrorData::SandboxUnreachable {
231 operation: operation.to_string(),
232 reason: "the Agent Platform API did not answer a sandbox read".to_string(),
233 })),
234 }
235 }
236
237 async fn await_operation(
242 &self,
243 operation: &str,
244 started: Operation,
245 ) -> Result<serde_json::Value> {
246 let Some(name) = started.name.clone() else {
247 return Err(AlienError::new(ErrorData::UnexpectedResponseFormat {
248 provider: "gcp-agent-platform".to_string(),
249 binding_name: operation.to_string(),
250 field: "name".to_string(),
251 response_json: "the operation carried no resource name to poll".to_string(),
252 }));
253 };
254
255 let mut current = started;
256 for _ in 0..OPERATION_POLL_ATTEMPTS {
257 if current.done == Some(true) {
258 return finish_operation(operation, &name, current);
259 }
260 tokio::time::sleep(OPERATION_POLL_INTERVAL).await;
261 current =
262 self.client
263 .get_operation(&name)
264 .await
265 .context(ErrorData::SandboxUnreachable {
266 operation: operation.to_string(),
267 reason: format!("could not read operation '{name}'"),
268 })?;
269 }
270
271 if current.done == Some(true) {
272 return finish_operation(operation, &name, current);
273 }
274 Err(AlienError::new(ErrorData::SandboxUnreachable {
275 operation: operation.to_string(),
276 reason: format!("operation '{name}' did not complete within its polling budget"),
277 }))
278 }
279
280 async fn execute_op(
286 &self,
287 sandbox_id: &str,
288 operation: &str,
289 envelope: serde_json::Value,
290 ) -> Result<Vec<u8>> {
291 let body = serde_json::to_vec(&envelope).map_err(|error| {
292 AlienError::new(ErrorData::SerializationFailed {
293 message: format!("could not encode the {operation} envelope: {error}"),
294 })
295 })?;
296
297 self.client
298 .execute(&self.engine, sandbox_id, &body)
299 .await
300 .map_err(|error| Self::execute_failed(operation, error))
301 }
302
303 fn execute_failed(
304 operation: &str,
305 error: AlienError<AgentPlatformErrorData>,
306 ) -> AlienError<ErrorData> {
307 if let Some(answer) = agent_answer(&error) {
308 return error.context(ErrorData::SandboxCommandFailed {
309 failure: "agentRefused".to_string(),
310 reason: format!("{operation} was refused: {answer}"),
311 });
312 }
313 if is_not_found(&error) {
314 return error.context(ErrorData::SandboxCommandFailed {
315 failure: "sandboxGone".to_string(),
316 reason: format!("{operation}: the sandbox does not exist"),
317 });
318 }
319 if operation == RUN_COMMAND || operation == JOB_START {
324 return error.context(ErrorData::SandboxOutcomeUnknown {
325 operation: operation.to_string(),
326 reason: "the sandbox did not complete the call".to_string(),
327 });
328 }
329 error.context(ErrorData::SandboxCommandFailed {
330 failure: "executeFailed".to_string(),
331 reason: format!("{operation} could not be completed against the sandbox"),
332 })
333 }
334
335 async fn probe_agent(&self, operation: &str, sandbox_id: &str) -> Result<u64> {
342 let unreachable = |reason: String| {
347 AlienError::new(ErrorData::SandboxUnreachable {
348 operation: operation.to_string(),
349 reason,
350 })
351 };
352
353 let body = tokio::time::timeout(
354 AGENT_PROBE_BUDGET,
355 self.client.execute(
356 &self.engine,
357 sandbox_id,
358 &serde_json::to_vec(&json!({ "v": AGENT_PROTOCOL_VERSION, "op": "health" }))
359 .unwrap_or_default(),
360 ),
361 )
362 .await
363 .map_err(|_| {
364 unreachable(format!(
365 "the sandbox's agent did not answer a health probe within {}s",
366 AGENT_PROBE_BUDGET.as_secs()
367 ))
368 })?
369 .map_err(|error| {
370 error.context(ErrorData::SandboxUnreachable {
371 operation: operation.to_string(),
372 reason: "the sandbox's agent did not answer a health probe".to_string(),
373 })
374 })?;
375
376 #[derive(Deserialize)]
377 #[serde(rename_all = "camelCase")]
378 struct Health {
379 protocol_version: u32,
380 boot_id: String,
381 }
382
383 let health: Health = serde_json::from_slice(&body).map_err(|_| {
384 unreachable(format!(
385 "the sandbox's agent answered a health probe with a body this provider cannot \
386 read: {}",
387 truncated(&body)
388 ))
389 })?;
390
391 if health.protocol_version != AGENT_PROTOCOL_VERSION {
392 return Err(unreachable(format!(
393 "the sandbox's agent speaks protocol {} where this provider speaks {}",
394 health.protocol_version, AGENT_PROTOCOL_VERSION
395 )));
396 }
397 if health.boot_id.is_empty() {
400 return Err(unreachable(
401 "the sandbox's agent reported no container boot id, so its identity cannot be \
402 established"
403 .to_string(),
404 ));
405 }
406 Ok(generation_from_boot_id(&health.boot_id))
407 }
408
409 async fn discard(
415 &self,
416 sandbox_id: &str,
417 reason: AlienError<ErrorData>,
418 ) -> AlienError<ErrorData> {
419 let Err(error) = self.client.delete_sandbox(&self.engine, sandbox_id).await else {
420 return reason;
421 };
422 warn!(
423 sandbox = %sandbox_id,
424 %error,
425 "could not delete a sandbox that was never handed to its caller"
426 );
427 reason.context(ErrorData::SandboxCommandFailed {
428 failure: "sandboxLeftBehind".to_string(),
429 reason: format!(
430 "sandbox '{sandbox_id}' was not handed to its caller and could not be deleted, so \
431 it is still running"
432 ),
433 })
434 }
435
436 async fn settle(&self, sandbox_id: &str) -> Result<u64> {
442 for _ in 0..SANDBOX_READY_ATTEMPTS {
443 let Some(sandbox) = self.read_sandbox(CREATE, sandbox_id).await? else {
444 return Err(AlienError::new(ErrorData::SandboxCommandFailed {
445 failure: "sandboxGone".to_string(),
446 reason: format!("sandbox '{sandbox_id}' disappeared while it was coming up"),
447 }));
448 };
449 match sandbox_state(CREATE, sandbox.state.as_deref())? {
450 SandboxState::Running => return self.wait_until_servable(sandbox_id).await,
451 SandboxState::Terminated => {
452 return Err(AlienError::new(ErrorData::SandboxCommandFailed {
453 failure: "sandboxTerminated".to_string(),
454 reason: format!(
455 "sandbox '{sandbox_id}' reached a terminal state while starting"
456 ),
457 }));
458 }
459 SandboxState::Starting | SandboxState::Paused => {}
463 }
464 tokio::time::sleep(SANDBOX_READY_INTERVAL).await;
465 }
466 Err(AlienError::new(ErrorData::SandboxUnreachable {
467 operation: CREATE.to_string(),
468 reason: format!(
469 "sandbox '{sandbox_id}' was not running after {}s",
470 SANDBOX_READY_ATTEMPTS as u64 * SANDBOX_READY_INTERVAL.as_secs()
471 ),
472 }))
473 }
474
475 async fn wait_until_servable(&self, sandbox_id: &str) -> Result<u64> {
479 let deadline = tokio::time::Instant::now() + AGENT_READY_TIMEOUT;
480 let mut last_error = None;
481 loop {
482 match tokio::time::timeout_at(deadline, self.probe_agent(CREATE, sandbox_id)).await {
484 Ok(Ok(generation)) => return Ok(generation),
485 Ok(Err(error)) => {
486 last_error = Some(error);
487 if tokio::time::Instant::now() + AGENT_READY_POLL >= deadline {
488 break;
489 }
490 tokio::time::sleep(AGENT_READY_POLL).await;
491 }
492 Err(_) => break,
493 }
494 }
495 let timed_out = ErrorData::SandboxUnreachable {
496 operation: CREATE.to_string(),
497 reason: format!(
498 "sandbox '{sandbox_id}' was running but its agent did not become servable within {}s",
499 AGENT_READY_TIMEOUT.as_secs()
500 ),
501 };
502 Err(match last_error {
503 Some(error) => error.context(timed_out),
504 None => AlienError::new(timed_out),
505 })
506 }
507
508 async fn run_synchronous(
510 &self,
511 sandbox_id: &str,
512 request: &RunCommandRequest,
513 ) -> Result<BoxStream<'static, Result<CommandOutput>>> {
514 let envelope = exec_envelope("exec", sandbox_id, request);
515 let body = self.execute_op(sandbox_id, RUN_COMMAND, envelope).await?;
516 let frames = parse_exec_frames(&body)?;
517 Ok(Box::pin(stream::iter(frames)))
518 }
519
520 async fn run_detached(
522 &self,
523 sandbox_id: &str,
524 request: RunCommandRequest,
525 ) -> Result<BoxStream<'static, Result<CommandOutput>>> {
526 let timeout = request.timeout;
527 let started = self.start_job(sandbox_id, request).await?;
528
529 let state = JobPollState {
530 client: self.client.clone(),
531 engine: self.engine.clone(),
532 sandbox_id: sandbox_id.to_string(),
533 job_id: started.job_id,
534 since_seq: None,
535 pending: VecDeque::new(),
536 finished: false,
537 stopped: false,
538 deadline_at: tokio::time::Instant::now() + timeout + JOB_POLL_GRACE,
539 };
540
541 Ok(Box::pin(stream::unfold(state, job_poll_step)))
542 }
543
544 fn checked_command(operation: &str, request: &RunCommandRequest) -> Result<()> {
546 if request.command.is_empty() {
547 return Err(AlienError::new(ErrorData::InvalidInput {
548 operation_context: operation.to_string(),
549 details: "a command must name a program to run".to_string(),
550 field_name: Some("command".to_string()),
551 }));
552 }
553 if timeout_millis(request.timeout) == 0 {
557 return Err(AlienError::new(ErrorData::SandboxCommandFailed {
558 failure: "invalidRequest".to_string(),
559 reason: "a command must carry a timeout of at least one millisecond".to_string(),
560 }));
561 }
562 Ok(())
563 }
564}
565
566impl Binding for GcpAgentPlatformSandbox {}
567
568#[async_trait]
569impl Sandbox for GcpAgentPlatformSandbox {
570 fn as_any(&self) -> &dyn std::any::Any {
571 self
572 }
573
574 fn capabilities(&self) -> SandboxCapabilities {
578 SandboxCapabilities::gcp_agent_platform()
579 }
580
581 async fn create(&self, request: CreateSandboxRequest) -> Result<SandboxInstance> {
582 if !request.env.is_empty() {
589 return Err(AlienError::new(ErrorData::OperationNotSupported {
590 operation: CREATE.to_string(),
591 reason: "Agent Platform sandboxes take no sandbox-level env; set env per command \
592 instead"
593 .to_string(),
594 }));
595 }
596
597 if request.tenant_key.is_some() {
600 return Err(AlienError::new(ErrorData::OperationNotSupported {
601 operation: CREATE.to_string(),
602 reason: "Agent Platform sandboxes take no tenantKey; create one sandbox per \
603 tenant instead"
604 .to_string(),
605 }));
606 }
607
608 let ttl = self
609 .lifetime_seconds(request.timeout_ms, CREATE)?
610 .map(|seconds| format!("{seconds}s"));
611
612 let started = self
613 .client
614 .create_sandbox(
615 &self.engine,
616 SandboxCreateRequest {
617 display_name: request.sandbox_id.clone(),
618 sandbox_environment_template: Some(self.template.clone()),
619 sandbox_environment_snapshot: None,
620 ttl,
621 },
622 )
623 .await
624 .context(ErrorData::SandboxUnreachable {
625 operation: CREATE.to_string(),
626 reason: "the Agent Platform API refused a sandbox create".to_string(),
627 })?;
628
629 let created: SandboxEnvironment = serde_json::from_value(
630 self.await_operation(CREATE, started).await?,
631 )
632 .map_err(|error| {
633 AlienError::new(ErrorData::UnexpectedResponseFormat {
634 provider: "gcp-agent-platform".to_string(),
635 binding_name: CREATE.to_string(),
636 field: "response".to_string(),
637 response_json: format!("the create operation resolved to a non-sandbox: {error}"),
638 })
639 })?;
640
641 let Some(sandbox_id) = created.name.as_deref().and_then(sandbox_segment) else {
646 return Err(AlienError::new(ErrorData::UnexpectedResponseFormat {
647 provider: "gcp-agent-platform".to_string(),
648 binding_name: CREATE.to_string(),
649 field: "name".to_string(),
650 response_json: format!("{:?}", created.name),
651 }));
652 };
653 let sandbox_id = sandbox_id.to_string();
654
655 match self.settle(&sandbox_id).await {
657 Ok(generation) => Ok(SandboxInstance {
658 sandbox_id,
659 state: SandboxState::Running,
660 generation,
661 }),
662 Err(error) => Err(self.discard(&sandbox_id, error).await),
663 }
664 }
665
666 async fn get(&self, sandbox_id: &str) -> Result<Option<SandboxInstance>> {
667 Self::checked_sandbox_id(GET, sandbox_id)?;
668 let Some(sandbox) = self.read_sandbox(GET, sandbox_id).await? else {
669 return Ok(None);
670 };
671
672 let state = sandbox_state(GET, sandbox.state.as_deref())?;
673 let generation = if state == SandboxState::Running {
677 self.probe_agent(GET, sandbox_id).await?
678 } else {
679 NO_GENERATION
680 };
681
682 Ok(Some(SandboxInstance {
683 sandbox_id: sandbox_id.to_string(),
684 state,
685 generation,
686 }))
687 }
688
689 async fn get_or_create(&self, request: CreateSandboxRequest) -> Result<ResolvedSandbox> {
690 if let Some(id) = request.sandbox_id.as_deref() {
691 match self.get(id).await {
695 Ok(Some(sandbox)) if sandbox.state == SandboxState::Running => {
696 return Ok(ResolvedSandbox::found(sandbox))
697 }
698 Ok(Some(sandbox)) if sandbox.state == SandboxState::Starting => {
702 let generation = self.settle(id).await?;
703 return Ok(ResolvedSandbox::found(SandboxInstance {
704 sandbox_id: id.to_string(),
705 state: SandboxState::Running,
706 generation,
707 }));
708 }
709 Ok(_) => {}
710 Err(error) if error.code == "SANDBOX_UNREACHABLE" => {}
711 Err(error) => {
712 return Err(error.context(ErrorData::SandboxCommandFailed {
713 failure: "getOrCreateFailed".to_string(),
714 reason: format!("{GET_OR_CREATE}: reaching sandbox '{id}' failed"),
715 }))
716 }
717 }
718 }
719
720 self.create(request).await.map(ResolvedSandbox::created)
721 }
722
723 async fn list(&self) -> Result<Vec<SandboxInstance>> {
724 let sandboxes = self.client.list_sandboxes(&self.engine).await.context(
725 ErrorData::SandboxUnreachable {
726 operation: "sandbox.list".to_string(),
727 reason: "the Agent Platform API did not answer a sandbox list".to_string(),
728 },
729 )?;
730
731 Ok(sandboxes
736 .into_iter()
737 .filter_map(|sandbox| {
738 let sandbox_id = sandbox.name.as_deref().and_then(sandbox_segment)?;
739 let state = sandbox_state("sandbox.list", sandbox.state.as_deref()).ok()?;
740 Some(SandboxInstance {
743 sandbox_id: sandbox_id.to_string(),
744 state,
745 generation: NO_GENERATION,
746 })
747 })
748 .collect())
749 }
750
751 async fn run_command(
752 &self,
753 sandbox_id: &str,
754 request: RunCommandRequest,
755 ) -> Result<BoxStream<'static, Result<CommandOutput>>> {
756 Self::checked_sandbox_id(RUN_COMMAND, sandbox_id)?;
757 Self::checked_command(RUN_COMMAND, &request)?;
758
759 if request.timeout <= MAX_SYNCHRONOUS_TIMEOUT {
762 self.run_synchronous(sandbox_id, &request).await
763 } else {
764 self.run_detached(sandbox_id, request).await
765 }
766 }
767
768 async fn start_job(&self, sandbox_id: &str, request: RunCommandRequest) -> Result<JobStart> {
769 Self::checked_sandbox_id(JOB_START, sandbox_id)?;
770 Self::checked_command(JOB_START, &request)?;
771
772 let envelope = exec_envelope("jobStart", sandbox_id, &request);
773 let body = self.execute_op(sandbox_id, JOB_START, envelope).await?;
774
775 let started: JobStartReply = serde_json::from_slice(&body).map_err(|_| {
779 AlienError::new(ErrorData::UnexpectedResponseFormat {
780 provider: "gcp-agent-platform".to_string(),
781 binding_name: JOB_START.to_string(),
782 field: "jobId".to_string(),
783 response_json: truncated(&body),
784 })
785 .context(ErrorData::SandboxOutcomeUnknown {
786 operation: JOB_START.to_string(),
787 reason: "the job started and its id could not be read, so it cannot be polled"
788 .to_string(),
789 })
790 })?;
791
792 Ok(JobStart {
793 job_id: started.job_id,
794 })
795 }
796
797 async fn poll_job(
798 &self,
799 sandbox_id: &str,
800 job_id: &str,
801 since_seq: Option<u64>,
802 ) -> Result<JobPoll> {
803 Self::checked_sandbox_id(JOB_POLL, sandbox_id)?;
804 let reply = poll_once(
805 self.client.as_ref(),
806 &self.engine,
807 sandbox_id,
808 job_id,
809 since_seq,
810 )
811 .await?;
812
813 Ok(JobPoll {
814 running: reply.running,
815 frames: reply
816 .frames
817 .into_iter()
818 .map(WireFrame::into_output)
819 .collect::<Result<Vec<_>>>()?,
820 exit: reply.exit_code.map(|code| JobExit {
821 code,
822 truncated: reply.truncated.unwrap_or(false),
823 }),
824 error: reply.error.map(|error| JobError {
825 code: error.code,
826 message: error.message,
827 }),
828 })
829 }
830
831 async fn cancel_job(&self, sandbox_id: &str, job_id: &str) -> Result<()> {
832 Self::checked_sandbox_id(JOB_CANCEL, sandbox_id)?;
833 let body = self
834 .client
835 .execute(&self.engine, sandbox_id, &cancel_body(job_id))
836 .await
837 .map_err(|error| unanswered_job(JOB_CANCEL, error))?;
838
839 if !cancel_confirmed(&body) {
840 return Err(AlienError::new(ErrorData::SandboxCommandFailed {
841 failure: "agentRefused".to_string(),
842 reason: format!("{JOB_CANCEL}: {}", truncated(&body)),
843 }));
844 }
845
846 Ok(())
847 }
848
849 async fn read_file(&self, sandbox_id: &str, path: &str) -> Result<Vec<u8>> {
850 Self::checked_sandbox_id("sandbox.readFile", sandbox_id)?;
851 let body = self
852 .execute_op(
853 sandbox_id,
854 "sandbox.readFile",
855 json!({ "v": AGENT_PROTOCOL_VERSION, "op": "readFile", "path": path }),
856 )
857 .await?;
858
859 #[derive(Deserialize)]
860 #[serde(rename_all = "camelCase")]
861 struct ReadFile {
862 contents_base64: String,
863 }
864 let read: ReadFile = serde_json::from_slice(&body).map_err(|_| {
865 AlienError::new(ErrorData::SandboxCommandFailed {
866 failure: "agentRefused".to_string(),
867 reason: format!("sandbox.readFile was refused: {}", truncated(&body)),
868 })
869 })?;
870
871 BASE64
872 .decode(read.contents_base64.as_bytes())
873 .map_err(|error| {
874 AlienError::new(ErrorData::UnexpectedResponseFormat {
875 provider: "gcp-agent-platform".to_string(),
876 binding_name: "sandbox.readFile".to_string(),
877 field: "contentsBase64".to_string(),
878 response_json: format!("the agent returned data that is not base64: {error}"),
879 })
880 })
881 }
882
883 async fn write_files(&self, sandbox_id: &str, files: BTreeMap<String, Vec<u8>>) -> Result<()> {
884 Self::checked_sandbox_id("sandbox.writeFiles", sandbox_id)?;
885 for (path, contents) in files {
889 let body = self
890 .execute_op(
891 sandbox_id,
892 "sandbox.writeFiles",
893 json!({
894 "v": AGENT_PROTOCOL_VERSION,
895 "op": "writeFile",
896 "path": path,
897 "contentsBase64": BASE64.encode(&contents),
898 }),
899 )
900 .await?;
901 confirm_empty_ok("sandbox.writeFiles", &body)?;
902 }
903 Ok(())
904 }
905
906 async fn preview(&self, _sandbox_id: &str, _port: u16) -> Result<PreviewCapability> {
907 Err(self.unsupported(
908 "preview",
909 "Agent Platform mints no port-scoped ingress capability; the only ingress is :execute",
910 ))
911 }
912
913 async fn pause(&self, sandbox_id: &str) -> Result<()> {
914 Self::checked_sandbox_id("sandbox.pause", sandbox_id)?;
915 Err(self.pause_resume_unsupported())
916 }
917
918 async fn resume(&self, sandbox_id: &str) -> Result<()> {
919 Self::checked_sandbox_id("sandbox.resume", sandbox_id)?;
920 Err(self.pause_resume_unsupported())
921 }
922
923 async fn snapshot(&self, sandbox_id: &str) -> Result<String> {
924 Self::checked_sandbox_id("sandbox.snapshot", sandbox_id)?;
925 let display_name = format!("snap-{}", uuid::Uuid::new_v4().simple());
929 let started = self
930 .client
931 .snapshot(&self.engine, sandbox_id, &display_name)
932 .await
933 .context(ErrorData::SandboxCommandFailed {
934 failure: "snapshotFailed".to_string(),
935 reason: format!("sandbox.snapshot: sandbox '{sandbox_id}' could not be captured"),
936 })?;
937
938 let snapshot: SandboxSnapshot =
939 serde_json::from_value(self.await_operation("sandbox.snapshot", started).await?)
940 .map_err(|error| {
941 AlienError::new(ErrorData::UnexpectedResponseFormat {
942 provider: "gcp-agent-platform".to_string(),
943 binding_name: "sandbox.snapshot".to_string(),
944 field: "response".to_string(),
945 response_json: format!(
946 "the snapshot operation resolved to a non-snapshot: {error}"
947 ),
948 })
949 })?;
950
951 snapshot.name.ok_or_else(|| {
952 AlienError::new(ErrorData::UnexpectedResponseFormat {
953 provider: "gcp-agent-platform".to_string(),
954 binding_name: "sandbox.snapshot".to_string(),
955 field: "name".to_string(),
956 response_json: "the snapshot completed without a resource name".to_string(),
957 })
958 })
959 }
960
961 async fn terminate(&self, sandbox_id: &str) -> Result<()> {
962 Self::checked_sandbox_id(TERMINATE, sandbox_id)?;
963 if let Err(error) = self.client.delete_sandbox(&self.engine, sandbox_id).await {
970 if !is_not_found(&error) {
971 return Err(error.context(ErrorData::SandboxUnreachable {
972 operation: TERMINATE.to_string(),
973 reason: format!("the delete of sandbox '{sandbox_id}' was not accepted"),
974 }));
975 }
976 return Ok(());
977 }
978
979 for _ in 0..TERMINATE_POLL_ATTEMPTS {
980 match self.client.get_sandbox(&self.engine, sandbox_id).await {
981 Err(error) if is_not_found(&error) => return Ok(()),
982 Err(error) => {
985 warn!(sandbox = %sandbox_id, %error, "could not confirm a sandbox is gone")
986 }
987 Ok(_) => {}
988 }
989 tokio::time::sleep(TERMINATE_POLL_INTERVAL).await;
990 }
991
992 Err(AlienError::new(ErrorData::SandboxUnreachable {
993 operation: TERMINATE.to_string(),
994 reason: format!(
995 "deletion of '{sandbox_id}' was accepted but the sandbox was still present after \
996 {}s; it may still be running",
997 TERMINATE_POLL_ATTEMPTS as u64 * TERMINATE_POLL_INTERVAL.as_secs()
998 ),
999 }))
1000 }
1001}
1002
1003async fn job_poll_step(mut state: JobPollState) -> Option<(Result<CommandOutput>, JobPollState)> {
1006 loop {
1007 if let Some(item) = state.pending.pop_front() {
1008 return Some((item, state));
1009 }
1010 if state.finished {
1011 return None;
1012 }
1013
1014 if tokio::time::Instant::now() >= state.deadline_at {
1015 let cancelled = state
1019 .client
1020 .execute(
1021 &state.engine,
1022 &state.sandbox_id,
1023 &cancel_body(&state.job_id),
1024 )
1025 .await;
1026 let confirmed = cancelled.as_ref().is_ok_and(|body| cancel_confirmed(body));
1027 state.pending.push_back(Err(match cancelled {
1028 Ok(_) if confirmed => AlienError::new(ErrorData::SandboxCommandFailed {
1029 failure: "timeoutExceeded".to_string(),
1030 reason: "the command's deadline elapsed before its job reported an outcome"
1031 .to_string(),
1032 }),
1033 Ok(_) => AlienError::new(ErrorData::SandboxOutcomeUnknown {
1034 operation: RUN_COMMAND.to_string(),
1035 reason: "the command's deadline elapsed and its job did not confirm the cancel"
1036 .to_string(),
1037 }),
1038 Err(error) => error.context(ErrorData::SandboxOutcomeUnknown {
1039 operation: RUN_COMMAND.to_string(),
1040 reason: "the command's deadline elapsed and its job could not be cancelled"
1041 .to_string(),
1042 }),
1043 }));
1044 state.finished = true;
1045 state.stopped = confirmed;
1048 continue;
1049 }
1050
1051 let poll = match poll_once(
1052 state.client.as_ref(),
1053 &state.engine,
1054 &state.sandbox_id,
1055 &state.job_id,
1056 state.since_seq,
1057 )
1058 .await
1059 {
1060 Ok(poll) => poll,
1061 Err(error) => {
1064 state
1065 .pending
1066 .push_back(Err(error.context(ErrorData::SandboxOutcomeUnknown {
1067 operation: RUN_COMMAND.to_string(),
1068 reason: "the job is no longer watched".to_string(),
1069 })));
1070 state.finished = true;
1071 continue;
1072 }
1073 };
1074
1075 for frame in poll.frames {
1076 state.since_seq = state.since_seq.max(frame.seq());
1080 let output = frame.into_output();
1081 let failed = output.is_err();
1084 state.pending.push_back(output);
1085 if failed {
1086 state.finished = true;
1087 break;
1088 }
1089 }
1090
1091 if !poll.running {
1092 let terminal = match poll.error {
1095 Some(error) => Err(AlienError::new(ErrorData::SandboxCommandFailed {
1096 failure: error.code,
1097 reason: error.message,
1098 })),
1099 None => match poll.exit_code {
1102 Some(code) => Ok(CommandOutput::Exit {
1103 code,
1104 truncated: poll.truncated.unwrap_or(false),
1105 }),
1106 None => Err(AlienError::new(ErrorData::SandboxOutcomeUnknown {
1107 operation: RUN_COMMAND.to_string(),
1108 reason: "the job finished without reporting an exit code".to_string(),
1109 })),
1110 },
1111 };
1112 state.pending.push_back(terminal);
1113 state.finished = true;
1114 state.stopped = true;
1115 continue;
1116 }
1117
1118 if state.pending.is_empty() {
1119 tokio::time::sleep(JOB_POLL_INTERVAL).await;
1120 }
1121 }
1122}
1123
1124async fn poll_once(
1127 client: &dyn AgentPlatformApi,
1128 engine: &str,
1129 sandbox_id: &str,
1130 job_id: &str,
1131 since_seq: Option<u64>,
1132) -> Result<JobPollReply> {
1133 let body = client
1134 .execute(engine, sandbox_id, &poll_body(job_id, since_seq))
1135 .await
1136 .map_err(|error| unanswered_job(JOB_POLL, error))?;
1137
1138 serde_json::from_slice(&body).map_err(|_| {
1139 AlienError::new(ErrorData::UnexpectedResponseFormat {
1140 provider: "gcp-agent-platform".to_string(),
1141 binding_name: JOB_POLL.to_string(),
1142 field: "jobPoll".to_string(),
1143 response_json: truncated(&body),
1144 })
1145 })
1146}
1147
1148fn unanswered_job(
1151 operation: &str,
1152 error: AlienError<AgentPlatformErrorData>,
1153) -> AlienError<ErrorData> {
1154 if is_not_found(&error) {
1155 return error.context(ErrorData::SandboxCommandFailed {
1156 failure: "sandboxGone".to_string(),
1157 reason: format!("{operation}: the sandbox does not exist"),
1158 });
1159 }
1160 error.context(ErrorData::SandboxUnreachable {
1161 operation: operation.to_string(),
1162 reason: "the sandbox did not complete the call".to_string(),
1163 })
1164}
1165
1166fn cancel_confirmed(body: &[u8]) -> bool {
1172 serde_json::from_slice::<serde_json::Value>(body)
1173 .is_ok_and(|value| value.as_object().is_some_and(serde_json::Map::is_empty))
1174}
1175
1176#[derive(Deserialize)]
1178#[serde(rename_all = "camelCase")]
1179struct JobStartReply {
1180 job_id: String,
1181}
1182
1183struct JobPollState {
1185 client: Arc<dyn AgentPlatformApi>,
1186 engine: String,
1187 sandbox_id: String,
1188 job_id: String,
1189 since_seq: Option<u64>,
1190 pending: VecDeque<Result<CommandOutput>>,
1191 finished: bool,
1192 stopped: bool,
1194 deadline_at: tokio::time::Instant,
1195}
1196
1197impl Drop for JobPollState {
1203 fn drop(&mut self) {
1204 if self.stopped {
1205 return;
1206 }
1207 let Ok(runtime) = tokio::runtime::Handle::try_current() else {
1208 warn!(
1209 sandbox = %self.sandbox_id,
1210 job = %self.job_id,
1211 "a job's stream was dropped outside a runtime, so the job runs to its timeout"
1212 );
1213 return;
1214 };
1215
1216 let client = self.client.clone();
1217 let engine = self.engine.clone();
1218 let sandbox_id = self.sandbox_id.clone();
1219 let job_id = self.job_id.clone();
1220 runtime.spawn(async move {
1221 match tokio::time::timeout(
1224 AGENT_PROBE_BUDGET,
1225 client.execute(&engine, &sandbox_id, &cancel_body(&job_id)),
1226 )
1227 .await
1228 {
1229 Ok(Ok(body)) if cancel_confirmed(&body) => {}
1230 Ok(Ok(body)) => warn!(
1231 sandbox = %sandbox_id,
1232 job = %job_id,
1233 reply = %truncated(&body),
1234 "a dropped job's cancel was refused, so the job runs to its timeout"
1235 ),
1236 Ok(Err(error)) => warn!(
1237 sandbox = %sandbox_id,
1238 job = %job_id,
1239 %error,
1240 "a dropped job's cancel did not land, so the job may run to its timeout"
1241 ),
1242 Err(_) => warn!(
1243 sandbox = %sandbox_id,
1244 job = %job_id,
1245 budget_secs = AGENT_PROBE_BUDGET.as_secs(),
1246 "a dropped job's cancel went unanswered, so the job may run to its timeout"
1247 ),
1248 }
1249 });
1250 }
1251}
1252
1253#[derive(Deserialize)]
1255#[serde(rename_all = "camelCase")]
1256struct JobPollReply {
1257 running: bool,
1258 #[serde(default)]
1259 frames: Vec<WireFrame>,
1260 #[serde(default)]
1261 exit_code: Option<i32>,
1262 #[serde(default)]
1263 truncated: Option<bool>,
1264 #[serde(default)]
1265 error: Option<JobErrorReply>,
1266}
1267
1268#[derive(Deserialize)]
1269#[serde(rename_all = "camelCase")]
1270struct JobErrorReply {
1271 code: String,
1272 message: String,
1273}
1274
1275#[derive(Deserialize)]
1277#[serde(rename_all = "camelCase", tag = "t")]
1278enum WireFrame {
1279 Stdout {
1280 seq: u64,
1281 data: String,
1282 },
1283 Stderr {
1284 seq: u64,
1285 data: String,
1286 },
1287 Exit {
1288 code: i32,
1289 #[serde(default)]
1290 truncated: bool,
1291 },
1292 Error {
1293 code: String,
1294 message: String,
1295 },
1296}
1297
1298impl WireFrame {
1299 fn is_terminal(&self) -> bool {
1300 matches!(self, Self::Exit { .. } | Self::Error { .. })
1301 }
1302
1303 fn seq(&self) -> Option<u64> {
1304 match self {
1305 Self::Stdout { seq, .. } | Self::Stderr { seq, .. } => Some(*seq),
1306 _ => None,
1307 }
1308 }
1309
1310 fn into_output(self) -> Result<CommandOutput> {
1311 match self {
1312 Self::Stdout { seq, data } => Ok(CommandOutput::Stdout {
1313 seq,
1314 data: decode_frame_data(&data)?,
1315 }),
1316 Self::Stderr { seq, data } => Ok(CommandOutput::Stderr {
1317 seq,
1318 data: decode_frame_data(&data)?,
1319 }),
1320 Self::Exit { code, truncated } => Ok(CommandOutput::Exit { code, truncated }),
1321 Self::Error { code, message } => {
1324 Err(AlienError::new(ErrorData::SandboxCommandFailed {
1325 failure: code,
1326 reason: message,
1327 }))
1328 }
1329 }
1330 }
1331}
1332
1333fn decode_frame_data(data: &str) -> Result<Vec<u8>> {
1336 BASE64.decode(data).map_err(|error| {
1337 AlienError::new(ErrorData::UnexpectedResponseFormat {
1338 provider: "gcp-agent-platform".to_string(),
1339 binding_name: RUN_COMMAND.to_string(),
1340 field: "data".to_string(),
1341 response_json: format!("an output frame's data is not base64: {error}"),
1342 })
1343 .context(ErrorData::SandboxOutcomeUnknown {
1344 operation: RUN_COMMAND.to_string(),
1345 reason: "an output frame did not decode".to_string(),
1346 })
1347 })
1348}
1349
1350fn parse_exec_frames(body: &[u8]) -> Result<Vec<Result<CommandOutput>>> {
1357 let mut frames = Vec::new();
1358 let mut saw_any = false;
1359 let mut saw_terminal = false;
1360
1361 for line in body.split(|byte| *byte == b'\n') {
1362 if line.is_empty() {
1363 continue;
1364 }
1365 match serde_json::from_slice::<WireFrame>(line) {
1366 Ok(frame) => {
1367 saw_any = true;
1368 saw_terminal |= frame.is_terminal();
1369 let output = frame.into_output();
1370 let failed = output.is_err();
1374 frames.push(output);
1375 if failed {
1376 saw_terminal = true;
1377 break;
1378 }
1379 }
1380 Err(error) => {
1381 if !saw_any {
1382 return Err(AlienError::new(ErrorData::SandboxCommandFailed {
1383 failure: "agentRefused".to_string(),
1384 reason: format!("run_command was refused: {}", truncated(body)),
1385 }));
1386 }
1387 frames.push(Err(AlienError::new(ErrorData::UnexpectedResponseFormat {
1390 provider: "gcp-agent-platform".to_string(),
1391 binding_name: RUN_COMMAND.to_string(),
1392 field: "frame".to_string(),
1393 response_json: format!("an output frame did not parse: {error}"),
1394 })
1395 .context(ErrorData::SandboxOutcomeUnknown {
1396 operation: RUN_COMMAND.to_string(),
1397 reason: "an output frame did not parse".to_string(),
1398 })));
1399 saw_terminal = true;
1400 break;
1401 }
1402 }
1403 }
1404
1405 if !saw_any {
1406 return Err(AlienError::new(ErrorData::SandboxCommandFailed {
1407 failure: "agentRefused".to_string(),
1408 reason: "run_command returned an empty body".to_string(),
1409 }));
1410 }
1411 if !saw_terminal {
1412 frames.push(Err(AlienError::new(ErrorData::SandboxOutcomeUnknown {
1413 operation: RUN_COMMAND.to_string(),
1414 reason: "the command's output ended without a terminal frame".to_string(),
1415 })));
1416 }
1417 Ok(frames)
1418}
1419
1420fn exec_envelope(op: &str, _sandbox_id: &str, request: &RunCommandRequest) -> serde_json::Value {
1423 json!({
1424 "v": AGENT_PROTOCOL_VERSION,
1425 "op": op,
1426 "command": request.argv(),
1427 "timeoutMs": timeout_millis(request.timeout),
1428 "cwd": request.cwd,
1429 "env": request.env,
1430 })
1431}
1432
1433fn poll_body(job_id: &str, since_seq: Option<u64>) -> Vec<u8> {
1434 serde_json::to_vec(&json!({
1435 "v": AGENT_PROTOCOL_VERSION,
1436 "op": "jobPoll",
1437 "jobId": job_id,
1438 "sinceSeq": since_seq,
1439 }))
1440 .unwrap_or_default()
1441}
1442
1443fn cancel_body(job_id: &str) -> Vec<u8> {
1444 serde_json::to_vec(&json!({
1445 "v": AGENT_PROTOCOL_VERSION,
1446 "op": "jobCancel",
1447 "jobId": job_id,
1448 }))
1449 .unwrap_or_default()
1450}
1451
1452fn timeout_millis(timeout: Duration) -> u64 {
1455 u64::try_from(timeout.as_millis()).unwrap_or(u64::MAX)
1456}
1457
1458fn confirm_empty_ok(operation: &str, body: &[u8]) -> Result<()> {
1463 if body.iter().all(|byte| byte.is_ascii_whitespace()) {
1464 return Ok(());
1465 }
1466 Err(AlienError::new(ErrorData::SandboxCommandFailed {
1467 failure: "agentRefused".to_string(),
1468 reason: format!("{operation} was refused: {}", truncated(body)),
1469 }))
1470}
1471
1472fn sandbox_segment(name: &str) -> Option<&str> {
1474 let segment = name.rsplit('/').next()?;
1475 is_addressable_id(segment).then_some(segment)
1476}
1477
1478fn is_addressable_id(id: &str) -> bool {
1479 !id.is_empty()
1480 && id.len() <= MAX_SANDBOX_ID
1481 && id
1482 .chars()
1483 .all(|c| c.is_ascii_alphanumeric() || c == '-' || c == '_')
1484}
1485
1486fn generation_from_boot_id(boot_id: &str) -> u64 {
1492 const FNV_OFFSET_BASIS: u64 = 0xcbf2_9ce4_8422_2325;
1493 const FNV_PRIME: u64 = 0x0000_0100_0000_01b3;
1494 let mut hash = FNV_OFFSET_BASIS;
1495 for byte in boot_id.as_bytes() {
1496 hash ^= u64::from(*byte);
1497 hash = hash.wrapping_mul(FNV_PRIME);
1498 }
1499 hash | 1
1500}
1501
1502fn sandbox_state(operation: &str, state: Option<&str>) -> Result<SandboxState> {
1505 match state {
1506 Some("STATE_RUNNING") => Ok(SandboxState::Running),
1507 Some("STATE_CREATING" | "STATE_PENDING" | "STATE_RESUMING") => Ok(SandboxState::Starting),
1508 Some("STATE_PAUSED" | "STATE_PAUSING" | "STATE_SUSPENDED") => Ok(SandboxState::Paused),
1509 Some("STATE_STOPPED" | "STATE_FAILED" | "STATE_DELETING" | "STATE_DELETED") => {
1510 Ok(SandboxState::Terminated)
1511 }
1512 other => Err(AlienError::new(ErrorData::UnexpectedResponseFormat {
1513 provider: "gcp-agent-platform".to_string(),
1514 binding_name: operation.to_string(),
1515 field: "state".to_string(),
1516 response_json: other
1517 .map_or_else(|| "absent".to_string(), |state| format!("\"{state}\"")),
1518 })),
1519 }
1520}
1521
1522fn finish_operation(operation: &str, name: &str, op: Operation) -> Result<serde_json::Value> {
1524 match op.result {
1525 Some(OperationResult::Response { response }) => Ok(response),
1526 Some(OperationResult::Error { error }) => {
1527 Err(AlienError::new(ErrorData::SandboxCommandFailed {
1528 failure: "operationFailed".to_string(),
1529 reason: format!(
1530 "{operation}: operation '{name}' failed (grpc {}): {}",
1531 error.code, error.message
1532 ),
1533 }))
1534 }
1535 None => Err(AlienError::new(ErrorData::UnexpectedResponseFormat {
1536 provider: "gcp-agent-platform".to_string(),
1537 binding_name: operation.to_string(),
1538 field: "response".to_string(),
1539 response_json: format!("operation '{name}' reported done without a result"),
1540 })),
1541 }
1542}
1543
1544fn agent_answer(error: &AlienError<AgentPlatformErrorData>) -> Option<String> {
1548 const RELAYED: &str = "Error Details: ";
1549 let message = super::refusal::captured_refusal(error)?;
1550 let details = message.split_once(RELAYED)?.1.trim();
1551 let code = details.split(':').next()?;
1552 let is_agent_code =
1553 !code.is_empty() && code.chars().all(|ch| ch.is_ascii_uppercase() || ch == '_');
1554 is_agent_code.then(|| truncated(details.as_bytes()))
1555}
1556
1557fn is_not_found(error: &AlienError<AgentPlatformErrorData>) -> bool {
1563 const NOT_FOUND: &str = "REMOTE_RESOURCE_NOT_FOUND";
1564 if error.code == NOT_FOUND {
1565 return true;
1566 }
1567 let mut node = error.source.as_deref();
1568 while let Some(current) = node {
1569 if current.code == NOT_FOUND {
1570 return true;
1571 }
1572 node = current.source.as_deref();
1573 }
1574 false
1575}
1576
1577fn truncated(body: &[u8]) -> String {
1579 const LIMIT: usize = 200;
1580 let text = String::from_utf8_lossy(body);
1581 let text = text.trim();
1582 if text.len() <= LIMIT {
1583 return text.to_string();
1584 }
1585 let end = (0..=LIMIT)
1586 .rev()
1587 .find(|at| text.is_char_boundary(*at))
1588 .unwrap_or(0);
1589 format!("{}…", &text[..end])
1590}
1591
1592#[cfg(test)]
1593#[path = "gcp_agent_platform_tests.rs"]
1594mod tests;