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, CreateSessionRequest, JobError, JobExit, JobPoll, JobStart,
24 PreviewCapability, RunCommandRequest, Sandbox, SandboxSession, SandboxSessionState,
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_DEADLINE: Duration = Duration::from_secs(30);
43
44const MAX_SESSION_ID: usize = 63;
47
48const SESSION_READY_ATTEMPTS: u32 = 150;
50const SESSION_READY_INTERVAL: Duration = Duration::from_secs(2);
51
52const OPERATION_POLL_ATTEMPTS: u32 = 150;
55const OPERATION_POLL_INTERVAL: Duration = Duration::from_secs(2);
56
57const TERMINATE_POLL_ATTEMPTS: u32 = 30;
60const TERMINATE_POLL_INTERVAL: Duration = Duration::from_secs(2);
61
62const JOB_POLL_INTERVAL: Duration = Duration::from_secs(1);
64
65const JOB_POLL_GRACE: Duration = Duration::from_secs(15);
69
70const CREATE: &str = "sandbox.create";
71const GET: &str = "sandbox.get";
72const GET_OR_CREATE: &str = "sandbox.getOrCreate";
73const RUN_COMMAND: &str = "sandbox.runCommand";
74const JOB_START: &str = "sandbox.jobStart";
75const JOB_POLL: &str = "sandbox.jobPoll";
76const JOB_CANCEL: &str = "sandbox.jobCancel";
77const TERMINATE: &str = "sandbox.terminate";
78
79const NO_GENERATION: u64 = 0;
84
85const AGENT_PROBE_BUDGET: Duration = Duration::from_secs(60);
92
93pub fn egress_control_config(
102 sandbox_label: &str,
103 egress: &SandboxEgress,
104) -> Result<EgressControlConfig> {
105 let Some(internet_access) = egress.internet_access_switch() else {
106 return Err(AlienError::new(ErrorData::InvalidInput {
107 operation_context: "sandbox.template".to_string(),
108 details: format!(
109 "sandbox '{sandbox_label}' asked for domain-scoped egress, which Agent \
110 Platform cannot express; it offers only 'allow' (open) and 'deny' (closed)"
111 ),
112 field_name: Some("egress".to_string()),
113 }));
114 };
115
116 Ok(EgressControlConfig {
117 internet_access: Some(internet_access),
118 extra: Default::default(),
119 })
120}
121
122#[derive(Debug)]
124pub struct GcpAgentPlatformSandbox {
125 client: Arc<dyn AgentPlatformApi>,
126 engine: String,
129 template: String,
131 session_ttl_seconds: Option<u32>,
133}
134
135impl GcpAgentPlatformSandbox {
136 pub fn new(
141 client: Arc<dyn AgentPlatformApi>,
142 engine: String,
143 template: String,
144 session_ttl_seconds: Option<u32>,
145 ) -> Self {
146 let engine = engine.rsplit('/').next().unwrap_or(&engine).to_string();
147 Self {
148 client,
149 engine,
150 template,
151 session_ttl_seconds,
152 }
153 }
154
155 #[cfg(test)]
158 pub(crate) fn engine(&self) -> &str {
159 &self.engine
160 }
161
162 fn unsupported(&self, capability: &str, reason: &str) -> AlienError<ErrorData> {
163 AlienError::new(ErrorData::OperationNotSupported {
164 operation: capability.to_string(),
165 reason: reason.to_string(),
166 })
167 }
168
169 fn checked_session_id(operation: &str, session_id: &str) -> Result<()> {
175 if is_addressable_id(session_id) {
176 return Ok(());
177 }
178 Err(AlienError::new(ErrorData::InvalidInput {
179 operation_context: operation.to_string(),
180 details: format!(
181 "session id '{session_id}' must be a single segment of letters, digits, '-' and \
182 '_', at most {MAX_SESSION_ID} characters"
183 ),
184 field_name: Some("sessionId".to_string()),
185 }))
186 }
187
188 async fn read_sandbox(
190 &self,
191 operation: &str,
192 session_id: &str,
193 ) -> Result<Option<SandboxEnvironment>> {
194 match self.client.get_sandbox(&self.engine, session_id).await {
195 Ok(sandbox) => Ok(Some(sandbox)),
196 Err(error) if is_not_found(&error) => Ok(None),
197 Err(error) => Err(error.context(ErrorData::SandboxUnreachable {
198 operation: operation.to_string(),
199 reason: "the Agent Platform API did not answer a sandbox read".to_string(),
200 })),
201 }
202 }
203
204 async fn await_operation(
209 &self,
210 operation: &str,
211 started: Operation,
212 ) -> Result<serde_json::Value> {
213 let Some(name) = started.name.clone() else {
214 return Err(AlienError::new(ErrorData::UnexpectedResponseFormat {
215 provider: "gcp-agent-platform".to_string(),
216 binding_name: operation.to_string(),
217 field: "name".to_string(),
218 response_json: "the operation carried no resource name to poll".to_string(),
219 }));
220 };
221
222 let mut current = started;
223 for _ in 0..OPERATION_POLL_ATTEMPTS {
224 if current.done == Some(true) {
225 return finish_operation(operation, &name, current);
226 }
227 tokio::time::sleep(OPERATION_POLL_INTERVAL).await;
228 current =
229 self.client
230 .get_operation(&name)
231 .await
232 .context(ErrorData::SandboxUnreachable {
233 operation: operation.to_string(),
234 reason: format!("could not read operation '{name}'"),
235 })?;
236 }
237
238 if current.done == Some(true) {
239 return finish_operation(operation, &name, current);
240 }
241 Err(AlienError::new(ErrorData::SandboxUnreachable {
242 operation: operation.to_string(),
243 reason: format!("operation '{name}' did not complete within its polling budget"),
244 }))
245 }
246
247 async fn execute_op(
253 &self,
254 session_id: &str,
255 operation: &str,
256 envelope: serde_json::Value,
257 ) -> Result<Vec<u8>> {
258 let body = serde_json::to_vec(&envelope).map_err(|error| {
259 AlienError::new(ErrorData::SerializationFailed {
260 message: format!("could not encode the {operation} envelope: {error}"),
261 })
262 })?;
263
264 self.client
265 .execute(&self.engine, session_id, &body)
266 .await
267 .map_err(|error| Self::execute_failed(operation, error))
268 }
269
270 fn execute_failed(
271 operation: &str,
272 error: AlienError<AgentPlatformErrorData>,
273 ) -> AlienError<ErrorData> {
274 if is_not_found(&error) {
275 return error.context(ErrorData::SandboxCommandFailed {
276 failure: "sessionGone".to_string(),
277 reason: format!("{operation}: the session does not exist"),
278 });
279 }
280 if operation == RUN_COMMAND || operation == JOB_START {
285 return error.context(ErrorData::SandboxOutcomeUnknown {
286 operation: operation.to_string(),
287 reason: "the session did not complete the call".to_string(),
288 });
289 }
290 error.context(ErrorData::SandboxCommandFailed {
291 failure: "executeFailed".to_string(),
292 reason: format!("{operation} could not be completed against the session"),
293 })
294 }
295
296 async fn probe_agent(&self, operation: &str, session_id: &str) -> Result<u64> {
303 let unreachable = |reason: String| {
308 AlienError::new(ErrorData::SandboxUnreachable {
309 operation: operation.to_string(),
310 reason,
311 })
312 };
313
314 let body = tokio::time::timeout(
315 AGENT_PROBE_BUDGET,
316 self.client.execute(
317 &self.engine,
318 session_id,
319 &serde_json::to_vec(&json!({ "v": AGENT_PROTOCOL_VERSION, "op": "health" }))
320 .unwrap_or_default(),
321 ),
322 )
323 .await
324 .map_err(|_| {
325 unreachable(format!(
326 "the session's agent did not answer a health probe within {}s",
327 AGENT_PROBE_BUDGET.as_secs()
328 ))
329 })?
330 .map_err(|error| {
331 error.context(ErrorData::SandboxUnreachable {
332 operation: operation.to_string(),
333 reason: "the session's agent did not answer a health probe".to_string(),
334 })
335 })?;
336
337 #[derive(Deserialize)]
338 #[serde(rename_all = "camelCase")]
339 struct Health {
340 protocol_version: u32,
341 boot_id: String,
342 }
343
344 let health: Health = serde_json::from_slice(&body).map_err(|_| {
345 unreachable(format!(
346 "the session's agent answered a health probe with a body this provider cannot \
347 read: {}",
348 truncated(&body)
349 ))
350 })?;
351
352 if health.protocol_version != AGENT_PROTOCOL_VERSION {
353 return Err(unreachable(format!(
354 "the session's agent speaks protocol {} where this provider speaks {}",
355 health.protocol_version, AGENT_PROTOCOL_VERSION
356 )));
357 }
358 if health.boot_id.is_empty() {
361 return Err(unreachable(
362 "the session's agent reported no container boot id, so its identity cannot be \
363 established"
364 .to_string(),
365 ));
366 }
367 Ok(generation_from_boot_id(&health.boot_id))
368 }
369
370 async fn discard(
376 &self,
377 session_id: &str,
378 reason: AlienError<ErrorData>,
379 ) -> AlienError<ErrorData> {
380 let Err(error) = self.client.delete_sandbox(&self.engine, session_id).await else {
381 return reason;
382 };
383 warn!(
384 session = %session_id,
385 %error,
386 "could not delete a sandbox that was never handed to its caller"
387 );
388 reason.context(ErrorData::SandboxCommandFailed {
389 failure: "sandboxLeftBehind".to_string(),
390 reason: format!(
391 "session '{session_id}' was not handed to its caller and could not be deleted, so \
392 it is still running"
393 ),
394 })
395 }
396
397 async fn settle(&self, session_id: &str) -> Result<u64> {
403 for _ in 0..SESSION_READY_ATTEMPTS {
404 let Some(sandbox) = self.read_sandbox(CREATE, session_id).await? else {
405 return Err(AlienError::new(ErrorData::SandboxCommandFailed {
406 failure: "sessionGone".to_string(),
407 reason: format!("session '{session_id}' disappeared while it was coming up"),
408 }));
409 };
410 match session_state(CREATE, sandbox.state.as_deref())? {
411 SandboxSessionState::Running => {
412 return self.probe_agent(CREATE, session_id).await;
413 }
414 SandboxSessionState::Terminated => {
415 return Err(AlienError::new(ErrorData::SandboxCommandFailed {
416 failure: "sessionTerminated".to_string(),
417 reason: format!(
418 "session '{session_id}' reached a terminal state while starting"
419 ),
420 }));
421 }
422 SandboxSessionState::Starting | SandboxSessionState::Suspended => {}
426 }
427 tokio::time::sleep(SESSION_READY_INTERVAL).await;
428 }
429 Err(AlienError::new(ErrorData::SandboxUnreachable {
430 operation: CREATE.to_string(),
431 reason: format!(
432 "session '{session_id}' was not running after {}s",
433 SESSION_READY_ATTEMPTS as u64 * SESSION_READY_INTERVAL.as_secs()
434 ),
435 }))
436 }
437
438 async fn run_synchronous(
440 &self,
441 session_id: &str,
442 request: &RunCommandRequest,
443 ) -> Result<BoxStream<'static, Result<CommandOutput>>> {
444 let envelope = exec_envelope("exec", session_id, request);
445 let body = self.execute_op(session_id, RUN_COMMAND, envelope).await?;
446 let frames = parse_exec_frames(&body)?;
447 Ok(Box::pin(stream::iter(frames)))
448 }
449
450 async fn run_detached(
452 &self,
453 session_id: &str,
454 request: RunCommandRequest,
455 ) -> Result<BoxStream<'static, Result<CommandOutput>>> {
456 let deadline = request.deadline;
457 let started = self.start_job(session_id, request).await?;
458
459 let state = JobPollState {
460 client: self.client.clone(),
461 engine: self.engine.clone(),
462 session_id: session_id.to_string(),
463 job_id: started.job_id,
464 since_seq: None,
465 pending: VecDeque::new(),
466 finished: false,
467 deadline_at: tokio::time::Instant::now() + deadline + JOB_POLL_GRACE,
468 };
469
470 Ok(Box::pin(stream::unfold(state, job_poll_step)))
471 }
472
473 fn checked_command(operation: &str, request: &RunCommandRequest) -> Result<()> {
475 if request.command.is_empty() {
476 return Err(AlienError::new(ErrorData::InvalidInput {
477 operation_context: operation.to_string(),
478 details: "a command must name a program to run".to_string(),
479 field_name: Some("command".to_string()),
480 }));
481 }
482 if deadline_millis(request.deadline) == 0 {
486 return Err(AlienError::new(ErrorData::SandboxCommandFailed {
487 failure: "invalidRequest".to_string(),
488 reason: "a command must carry a deadline of at least one millisecond".to_string(),
489 }));
490 }
491 Ok(())
492 }
493}
494
495impl Binding for GcpAgentPlatformSandbox {}
496
497#[async_trait]
498impl Sandbox for GcpAgentPlatformSandbox {
499 fn as_any(&self) -> &dyn std::any::Any {
500 self
501 }
502
503 fn capabilities(&self) -> SandboxCapabilities {
507 SandboxCapabilities::gcp_agent_platform()
508 }
509
510 async fn create(&self, request: CreateSessionRequest) -> Result<SandboxSession> {
511 if !request.env.is_empty() {
518 return Err(AlienError::new(ErrorData::OperationNotSupported {
519 operation: CREATE.to_string(),
520 reason: "Agent Platform sandboxes take no session-level env; set env per command \
521 instead"
522 .to_string(),
523 }));
524 }
525
526 if request.tenant_key.is_some() {
529 return Err(AlienError::new(ErrorData::OperationNotSupported {
530 operation: CREATE.to_string(),
531 reason: "Agent Platform sandboxes take no tenantKey; create one sandbox per \
532 tenant instead"
533 .to_string(),
534 }));
535 }
536
537 let started = self
538 .client
539 .create_sandbox(
540 &self.engine,
541 SandboxCreateRequest {
542 display_name: request.session_id.clone(),
543 sandbox_environment_template: Some(self.template.clone()),
544 sandbox_environment_snapshot: None,
545 ttl: self
546 .session_ttl_seconds
547 .map(|seconds| format!("{seconds}s")),
548 },
549 )
550 .await
551 .context(ErrorData::SandboxUnreachable {
552 operation: CREATE.to_string(),
553 reason: "the Agent Platform API refused a sandbox create".to_string(),
554 })?;
555
556 let created: SandboxEnvironment = serde_json::from_value(
557 self.await_operation(CREATE, started).await?,
558 )
559 .map_err(|error| {
560 AlienError::new(ErrorData::UnexpectedResponseFormat {
561 provider: "gcp-agent-platform".to_string(),
562 binding_name: CREATE.to_string(),
563 field: "response".to_string(),
564 response_json: format!("the create operation resolved to a non-sandbox: {error}"),
565 })
566 })?;
567
568 let Some(session_id) = created.name.as_deref().and_then(session_segment) else {
573 return Err(AlienError::new(ErrorData::UnexpectedResponseFormat {
574 provider: "gcp-agent-platform".to_string(),
575 binding_name: CREATE.to_string(),
576 field: "name".to_string(),
577 response_json: format!("{:?}", created.name),
578 }));
579 };
580 let session_id = session_id.to_string();
581
582 match self.settle(&session_id).await {
584 Ok(generation) => Ok(SandboxSession {
585 session_id,
586 state: SandboxSessionState::Running,
587 generation,
588 }),
589 Err(error) => Err(self.discard(&session_id, error).await),
590 }
591 }
592
593 async fn get(&self, session_id: &str) -> Result<Option<SandboxSession>> {
594 Self::checked_session_id(GET, session_id)?;
595 let Some(sandbox) = self.read_sandbox(GET, session_id).await? else {
596 return Ok(None);
597 };
598
599 let state = session_state(GET, sandbox.state.as_deref())?;
600 let generation = if state == SandboxSessionState::Running {
604 self.probe_agent(GET, session_id).await?
605 } else {
606 NO_GENERATION
607 };
608
609 Ok(Some(SandboxSession {
610 session_id: session_id.to_string(),
611 state,
612 generation,
613 }))
614 }
615
616 async fn get_or_create(&self, request: CreateSessionRequest) -> Result<SandboxSession> {
617 if let Some(id) = request.session_id.as_deref() {
618 match self.get(id).await {
622 Ok(Some(session)) if session.state == SandboxSessionState::Running => {
623 return Ok(session)
624 }
625 Ok(Some(session)) if session.state == SandboxSessionState::Suspended => {
631 if self.resume(id).await.is_ok() {
632 match self.get(id).await {
633 Ok(Some(woken)) if woken.state == SandboxSessionState::Running => {
634 return Ok(woken)
635 }
636 _ => {
637 if let Err(error) = self.suspend(id).await {
641 return Err(error.context(ErrorData::SandboxCommandFailed {
642 failure: "resumeRollbackFailed".to_string(),
643 reason: format!(
644 "{GET_OR_CREATE}: woke session '{id}' but could not \
645 confirm it healthy or put it back to sleep"
646 ),
647 }));
648 }
649 }
650 }
651 }
652 }
653 Ok(_) => {}
654 Err(error) if error.code == "SANDBOX_UNREACHABLE" => {}
655 Err(error) => {
656 return Err(error.context(ErrorData::SandboxCommandFailed {
657 failure: "getOrCreateFailed".to_string(),
658 reason: format!("{GET_OR_CREATE}: reaching session '{id}' failed"),
659 }))
660 }
661 }
662 }
663
664 self.create(request).await
665 }
666
667 async fn list(&self) -> Result<Vec<SandboxSession>> {
668 let sandboxes = self.client.list_sandboxes(&self.engine).await.context(
669 ErrorData::SandboxUnreachable {
670 operation: "sandbox.list".to_string(),
671 reason: "the Agent Platform API did not answer a sandbox list".to_string(),
672 },
673 )?;
674
675 Ok(sandboxes
680 .into_iter()
681 .filter_map(|sandbox| {
682 let session_id = sandbox.name.as_deref().and_then(session_segment)?;
683 let state = session_state("sandbox.list", sandbox.state.as_deref()).ok()?;
684 Some(SandboxSession {
687 session_id: session_id.to_string(),
688 state,
689 generation: NO_GENERATION,
690 })
691 })
692 .collect())
693 }
694
695 async fn run_command(
696 &self,
697 session_id: &str,
698 request: RunCommandRequest,
699 ) -> Result<BoxStream<'static, Result<CommandOutput>>> {
700 Self::checked_session_id(RUN_COMMAND, session_id)?;
701 Self::checked_command(RUN_COMMAND, &request)?;
702
703 if request.deadline <= MAX_SYNCHRONOUS_DEADLINE {
706 self.run_synchronous(session_id, &request).await
707 } else {
708 self.run_detached(session_id, request).await
709 }
710 }
711
712 async fn start_job(&self, session_id: &str, request: RunCommandRequest) -> Result<JobStart> {
713 Self::checked_session_id(JOB_START, session_id)?;
714 Self::checked_command(JOB_START, &request)?;
715
716 let envelope = exec_envelope("jobStart", session_id, &request);
717 let body = self.execute_op(session_id, JOB_START, envelope).await?;
718
719 let started: JobStartReply = serde_json::from_slice(&body).map_err(|_| {
723 AlienError::new(ErrorData::UnexpectedResponseFormat {
724 provider: "gcp-agent-platform".to_string(),
725 binding_name: JOB_START.to_string(),
726 field: "jobId".to_string(),
727 response_json: truncated(&body),
728 })
729 .context(ErrorData::SandboxOutcomeUnknown {
730 operation: JOB_START.to_string(),
731 reason: "the job started and its id could not be read, so it cannot be polled"
732 .to_string(),
733 })
734 })?;
735
736 Ok(JobStart {
737 job_id: started.job_id,
738 })
739 }
740
741 async fn poll_job(
742 &self,
743 session_id: &str,
744 job_id: &str,
745 since_seq: Option<u64>,
746 ) -> Result<JobPoll> {
747 Self::checked_session_id(JOB_POLL, session_id)?;
748 let reply = poll_once(
749 self.client.as_ref(),
750 &self.engine,
751 session_id,
752 job_id,
753 since_seq,
754 )
755 .await?;
756
757 Ok(JobPoll {
758 running: reply.running,
759 frames: reply
760 .frames
761 .into_iter()
762 .map(WireFrame::into_output)
763 .collect::<Result<Vec<_>>>()?,
764 exit: reply.exit_code.map(|code| JobExit {
765 code,
766 truncated: reply.truncated.unwrap_or(false),
767 }),
768 error: reply.error.map(|error| JobError {
769 code: error.code,
770 message: error.message,
771 }),
772 })
773 }
774
775 async fn cancel_job(&self, session_id: &str, job_id: &str) -> Result<()> {
776 Self::checked_session_id(JOB_CANCEL, session_id)?;
777 let body = self
778 .client
779 .execute(&self.engine, session_id, &cancel_body(job_id))
780 .await
781 .map_err(|error| unanswered_job(JOB_CANCEL, error))?;
782
783 if !cancel_confirmed(&body) {
784 return Err(AlienError::new(ErrorData::SandboxCommandFailed {
785 failure: "agentRefused".to_string(),
786 reason: format!("{JOB_CANCEL}: {}", truncated(&body)),
787 }));
788 }
789
790 Ok(())
791 }
792
793 async fn read_file(&self, session_id: &str, path: &str) -> Result<Vec<u8>> {
794 Self::checked_session_id("sandbox.readFile", session_id)?;
795 let body = self
796 .execute_op(
797 session_id,
798 "sandbox.readFile",
799 json!({ "v": AGENT_PROTOCOL_VERSION, "op": "readFile", "path": path }),
800 )
801 .await?;
802
803 #[derive(Deserialize)]
804 #[serde(rename_all = "camelCase")]
805 struct ReadFile {
806 contents_base64: String,
807 }
808 let read: ReadFile = serde_json::from_slice(&body).map_err(|_| {
809 AlienError::new(ErrorData::SandboxCommandFailed {
810 failure: "agentRefused".to_string(),
811 reason: format!("sandbox.readFile was refused: {}", truncated(&body)),
812 })
813 })?;
814
815 BASE64
816 .decode(read.contents_base64.as_bytes())
817 .map_err(|error| {
818 AlienError::new(ErrorData::UnexpectedResponseFormat {
819 provider: "gcp-agent-platform".to_string(),
820 binding_name: "sandbox.readFile".to_string(),
821 field: "contentsBase64".to_string(),
822 response_json: format!("the agent returned data that is not base64: {error}"),
823 })
824 })
825 }
826
827 async fn write_files(&self, session_id: &str, files: BTreeMap<String, Vec<u8>>) -> Result<()> {
828 Self::checked_session_id("sandbox.writeFiles", session_id)?;
829 for (path, contents) in files {
833 let body = self
834 .execute_op(
835 session_id,
836 "sandbox.writeFiles",
837 json!({
838 "v": AGENT_PROTOCOL_VERSION,
839 "op": "writeFile",
840 "path": path,
841 "contentsBase64": BASE64.encode(&contents),
842 }),
843 )
844 .await?;
845 confirm_empty_ok("sandbox.writeFiles", &body)?;
846 }
847 Ok(())
848 }
849
850 async fn mkdir(&self, session_id: &str, path: &str) -> Result<()> {
851 Self::checked_session_id("sandbox.mkdir", session_id)?;
852 let body = self
853 .execute_op(
854 session_id,
855 "sandbox.mkdir",
856 json!({ "v": AGENT_PROTOCOL_VERSION, "op": "mkdir", "path": path }),
857 )
858 .await?;
859 confirm_empty_ok("sandbox.mkdir", &body)
860 }
861
862 async fn preview(&self, _session_id: &str, _port: u16) -> Result<PreviewCapability> {
863 Err(self.unsupported(
864 "preview",
865 "Agent Platform mints no port-scoped ingress capability; the only ingress is :execute",
866 ))
867 }
868
869 async fn suspend(&self, session_id: &str) -> Result<()> {
870 Self::checked_session_id("sandbox.suspend", session_id)?;
871 let started = self.client.pause(&self.engine, session_id).await.context(
872 ErrorData::SandboxCommandFailed {
873 failure: "suspendFailed".to_string(),
874 reason: format!("sandbox.suspend: session '{session_id}' could not be paused"),
875 },
876 )?;
877 self.await_operation("sandbox.suspend", started).await?;
878 Ok(())
879 }
880
881 async fn resume(&self, session_id: &str) -> Result<()> {
882 Self::checked_session_id("sandbox.resume", session_id)?;
883 let started = self.client.resume(&self.engine, session_id).await.context(
884 ErrorData::SandboxCommandFailed {
885 failure: "resumeFailed".to_string(),
886 reason: format!("sandbox.resume: session '{session_id}' could not be resumed"),
887 },
888 )?;
889 self.await_operation("sandbox.resume", started).await?;
890 Ok(())
891 }
892
893 async fn snapshot(&self, session_id: &str) -> Result<String> {
894 Self::checked_session_id("sandbox.snapshot", session_id)?;
895 let display_name = format!("snap-{}", uuid::Uuid::new_v4().simple());
899 let started = self
900 .client
901 .snapshot(&self.engine, session_id, &display_name)
902 .await
903 .context(ErrorData::SandboxCommandFailed {
904 failure: "snapshotFailed".to_string(),
905 reason: format!("sandbox.snapshot: session '{session_id}' could not be captured"),
906 })?;
907
908 let snapshot: SandboxSnapshot =
909 serde_json::from_value(self.await_operation("sandbox.snapshot", started).await?)
910 .map_err(|error| {
911 AlienError::new(ErrorData::UnexpectedResponseFormat {
912 provider: "gcp-agent-platform".to_string(),
913 binding_name: "sandbox.snapshot".to_string(),
914 field: "response".to_string(),
915 response_json: format!(
916 "the snapshot operation resolved to a non-snapshot: {error}"
917 ),
918 })
919 })?;
920
921 snapshot.name.ok_or_else(|| {
922 AlienError::new(ErrorData::UnexpectedResponseFormat {
923 provider: "gcp-agent-platform".to_string(),
924 binding_name: "sandbox.snapshot".to_string(),
925 field: "name".to_string(),
926 response_json: "the snapshot completed without a resource name".to_string(),
927 })
928 })
929 }
930
931 async fn terminate(&self, session_id: &str) -> Result<()> {
932 Self::checked_session_id(TERMINATE, session_id)?;
933 if let Err(error) = self.client.delete_sandbox(&self.engine, session_id).await {
940 if !is_not_found(&error) {
941 return Err(error.context(ErrorData::SandboxUnreachable {
942 operation: TERMINATE.to_string(),
943 reason: format!("the delete of session '{session_id}' was not accepted"),
944 }));
945 }
946 return Ok(());
947 }
948
949 for _ in 0..TERMINATE_POLL_ATTEMPTS {
950 match self.client.get_sandbox(&self.engine, session_id).await {
951 Err(error) if is_not_found(&error) => return Ok(()),
952 Err(error) => {
955 warn!(session = %session_id, %error, "could not confirm a sandbox is gone")
956 }
957 Ok(_) => {}
958 }
959 tokio::time::sleep(TERMINATE_POLL_INTERVAL).await;
960 }
961
962 Err(AlienError::new(ErrorData::SandboxUnreachable {
963 operation: TERMINATE.to_string(),
964 reason: format!(
965 "deletion of '{session_id}' was accepted but the session was still present after \
966 {}s; it may still be running",
967 TERMINATE_POLL_ATTEMPTS as u64 * TERMINATE_POLL_INTERVAL.as_secs()
968 ),
969 }))
970 }
971}
972
973async fn job_poll_step(mut state: JobPollState) -> Option<(Result<CommandOutput>, JobPollState)> {
976 loop {
977 if let Some(item) = state.pending.pop_front() {
978 return Some((item, state));
979 }
980 if state.finished {
981 return None;
982 }
983
984 if tokio::time::Instant::now() >= state.deadline_at {
985 let cancelled = state
989 .client
990 .execute(
991 &state.engine,
992 &state.session_id,
993 &cancel_body(&state.job_id),
994 )
995 .await;
996 let confirmed = cancelled.as_ref().is_ok_and(|body| cancel_confirmed(body));
997 state.pending.push_back(Err(match cancelled {
998 Ok(_) if confirmed => AlienError::new(ErrorData::SandboxCommandFailed {
999 failure: "deadlineExceeded".to_string(),
1000 reason: "the command's deadline elapsed before its job reported an outcome"
1001 .to_string(),
1002 }),
1003 Ok(_) => AlienError::new(ErrorData::SandboxOutcomeUnknown {
1004 operation: RUN_COMMAND.to_string(),
1005 reason: "the command's deadline elapsed and its job did not confirm the cancel"
1006 .to_string(),
1007 }),
1008 Err(error) => error.context(ErrorData::SandboxOutcomeUnknown {
1009 operation: RUN_COMMAND.to_string(),
1010 reason: "the command's deadline elapsed and its job could not be cancelled"
1011 .to_string(),
1012 }),
1013 }));
1014 state.finished = true;
1015 continue;
1016 }
1017
1018 let poll = match poll_once(
1019 state.client.as_ref(),
1020 &state.engine,
1021 &state.session_id,
1022 &state.job_id,
1023 state.since_seq,
1024 )
1025 .await
1026 {
1027 Ok(poll) => poll,
1028 Err(error) => {
1031 state
1032 .pending
1033 .push_back(Err(error.context(ErrorData::SandboxOutcomeUnknown {
1034 operation: RUN_COMMAND.to_string(),
1035 reason: "the job is no longer watched".to_string(),
1036 })));
1037 state.finished = true;
1038 continue;
1039 }
1040 };
1041
1042 for frame in poll.frames {
1043 state.since_seq = state.since_seq.max(frame.seq());
1047 let output = frame.into_output();
1048 let failed = output.is_err();
1051 state.pending.push_back(output);
1052 if failed {
1053 state.finished = true;
1054 break;
1055 }
1056 }
1057
1058 if !poll.running {
1059 let terminal = match poll.error {
1062 Some(error) => Err(AlienError::new(ErrorData::SandboxCommandFailed {
1063 failure: error.code,
1064 reason: error.message,
1065 })),
1066 None => match poll.exit_code {
1069 Some(code) => Ok(CommandOutput::Exit {
1070 code,
1071 truncated: poll.truncated.unwrap_or(false),
1072 }),
1073 None => Err(AlienError::new(ErrorData::SandboxOutcomeUnknown {
1074 operation: RUN_COMMAND.to_string(),
1075 reason: "the job finished without reporting an exit code".to_string(),
1076 })),
1077 },
1078 };
1079 state.pending.push_back(terminal);
1080 state.finished = true;
1081 continue;
1082 }
1083
1084 if state.pending.is_empty() {
1085 tokio::time::sleep(JOB_POLL_INTERVAL).await;
1086 }
1087 }
1088}
1089
1090async fn poll_once(
1093 client: &dyn AgentPlatformApi,
1094 engine: &str,
1095 session_id: &str,
1096 job_id: &str,
1097 since_seq: Option<u64>,
1098) -> Result<JobPollReply> {
1099 let body = client
1100 .execute(engine, session_id, &poll_body(job_id, since_seq))
1101 .await
1102 .map_err(|error| unanswered_job(JOB_POLL, error))?;
1103
1104 serde_json::from_slice(&body).map_err(|_| {
1105 AlienError::new(ErrorData::UnexpectedResponseFormat {
1106 provider: "gcp-agent-platform".to_string(),
1107 binding_name: JOB_POLL.to_string(),
1108 field: "jobPoll".to_string(),
1109 response_json: truncated(&body),
1110 })
1111 })
1112}
1113
1114fn unanswered_job(
1117 operation: &str,
1118 error: AlienError<AgentPlatformErrorData>,
1119) -> AlienError<ErrorData> {
1120 if is_not_found(&error) {
1121 return error.context(ErrorData::SandboxCommandFailed {
1122 failure: "sessionGone".to_string(),
1123 reason: format!("{operation}: the session does not exist"),
1124 });
1125 }
1126 error.context(ErrorData::SandboxUnreachable {
1127 operation: operation.to_string(),
1128 reason: "the session did not complete the call".to_string(),
1129 })
1130}
1131
1132fn cancel_confirmed(body: &[u8]) -> bool {
1138 serde_json::from_slice::<serde_json::Value>(body).is_ok_and(|value| value.is_object())
1139}
1140
1141#[derive(Deserialize)]
1143#[serde(rename_all = "camelCase")]
1144struct JobStartReply {
1145 job_id: String,
1146}
1147
1148struct JobPollState {
1150 client: Arc<dyn AgentPlatformApi>,
1151 engine: String,
1152 session_id: String,
1153 job_id: String,
1154 since_seq: Option<u64>,
1155 pending: VecDeque<Result<CommandOutput>>,
1156 finished: bool,
1157 deadline_at: tokio::time::Instant,
1158}
1159
1160#[derive(Deserialize)]
1162#[serde(rename_all = "camelCase")]
1163struct JobPollReply {
1164 running: bool,
1165 #[serde(default)]
1166 frames: Vec<WireFrame>,
1167 #[serde(default)]
1168 exit_code: Option<i32>,
1169 #[serde(default)]
1170 truncated: Option<bool>,
1171 #[serde(default)]
1172 error: Option<JobErrorReply>,
1173}
1174
1175#[derive(Deserialize)]
1176#[serde(rename_all = "camelCase")]
1177struct JobErrorReply {
1178 code: String,
1179 message: String,
1180}
1181
1182#[derive(Deserialize)]
1184#[serde(rename_all = "camelCase", tag = "t")]
1185enum WireFrame {
1186 Stdout {
1187 seq: u64,
1188 data: String,
1189 },
1190 Stderr {
1191 seq: u64,
1192 data: String,
1193 },
1194 Exit {
1195 code: i32,
1196 #[serde(default)]
1197 truncated: bool,
1198 },
1199 Error {
1200 code: String,
1201 message: String,
1202 },
1203}
1204
1205impl WireFrame {
1206 fn is_terminal(&self) -> bool {
1207 matches!(self, Self::Exit { .. } | Self::Error { .. })
1208 }
1209
1210 fn seq(&self) -> Option<u64> {
1211 match self {
1212 Self::Stdout { seq, .. } | Self::Stderr { seq, .. } => Some(*seq),
1213 _ => None,
1214 }
1215 }
1216
1217 fn into_output(self) -> Result<CommandOutput> {
1218 match self {
1219 Self::Stdout { seq, data } => Ok(CommandOutput::Stdout {
1220 seq,
1221 data: decode_frame_data(&data)?,
1222 }),
1223 Self::Stderr { seq, data } => Ok(CommandOutput::Stderr {
1224 seq,
1225 data: decode_frame_data(&data)?,
1226 }),
1227 Self::Exit { code, truncated } => Ok(CommandOutput::Exit { code, truncated }),
1228 Self::Error { code, message } => {
1231 Err(AlienError::new(ErrorData::SandboxCommandFailed {
1232 failure: code,
1233 reason: message,
1234 }))
1235 }
1236 }
1237 }
1238}
1239
1240fn decode_frame_data(data: &str) -> Result<Vec<u8>> {
1243 BASE64.decode(data).map_err(|error| {
1244 AlienError::new(ErrorData::UnexpectedResponseFormat {
1245 provider: "gcp-agent-platform".to_string(),
1246 binding_name: RUN_COMMAND.to_string(),
1247 field: "data".to_string(),
1248 response_json: format!("an output frame's data is not base64: {error}"),
1249 })
1250 .context(ErrorData::SandboxOutcomeUnknown {
1251 operation: RUN_COMMAND.to_string(),
1252 reason: "an output frame did not decode".to_string(),
1253 })
1254 })
1255}
1256
1257fn parse_exec_frames(body: &[u8]) -> Result<Vec<Result<CommandOutput>>> {
1264 let mut frames = Vec::new();
1265 let mut saw_any = false;
1266 let mut saw_terminal = false;
1267
1268 for line in body.split(|byte| *byte == b'\n') {
1269 if line.is_empty() {
1270 continue;
1271 }
1272 match serde_json::from_slice::<WireFrame>(line) {
1273 Ok(frame) => {
1274 saw_any = true;
1275 saw_terminal |= frame.is_terminal();
1276 let output = frame.into_output();
1277 let failed = output.is_err();
1281 frames.push(output);
1282 if failed {
1283 saw_terminal = true;
1284 break;
1285 }
1286 }
1287 Err(error) => {
1288 if !saw_any {
1289 return Err(AlienError::new(ErrorData::SandboxCommandFailed {
1290 failure: "agentRefused".to_string(),
1291 reason: format!("run_command was refused: {}", truncated(body)),
1292 }));
1293 }
1294 frames.push(Err(AlienError::new(ErrorData::UnexpectedResponseFormat {
1297 provider: "gcp-agent-platform".to_string(),
1298 binding_name: RUN_COMMAND.to_string(),
1299 field: "frame".to_string(),
1300 response_json: format!("an output frame did not parse: {error}"),
1301 })
1302 .context(ErrorData::SandboxOutcomeUnknown {
1303 operation: RUN_COMMAND.to_string(),
1304 reason: "an output frame did not parse".to_string(),
1305 })));
1306 saw_terminal = true;
1307 break;
1308 }
1309 }
1310 }
1311
1312 if !saw_any {
1313 return Err(AlienError::new(ErrorData::SandboxCommandFailed {
1314 failure: "agentRefused".to_string(),
1315 reason: "run_command returned an empty body".to_string(),
1316 }));
1317 }
1318 if !saw_terminal {
1319 frames.push(Err(AlienError::new(ErrorData::SandboxOutcomeUnknown {
1320 operation: RUN_COMMAND.to_string(),
1321 reason: "the command's output ended without a terminal frame".to_string(),
1322 })));
1323 }
1324 Ok(frames)
1325}
1326
1327fn exec_envelope(op: &str, _session_id: &str, request: &RunCommandRequest) -> serde_json::Value {
1330 json!({
1331 "v": AGENT_PROTOCOL_VERSION,
1332 "op": op,
1333 "command": request.command,
1334 "deadlineMs": deadline_millis(request.deadline),
1335 "workingDirectory": request.working_directory,
1336 "env": request.env,
1337 })
1338}
1339
1340fn poll_body(job_id: &str, since_seq: Option<u64>) -> Vec<u8> {
1341 serde_json::to_vec(&json!({
1342 "v": AGENT_PROTOCOL_VERSION,
1343 "op": "jobPoll",
1344 "jobId": job_id,
1345 "sinceSeq": since_seq,
1346 }))
1347 .unwrap_or_default()
1348}
1349
1350fn cancel_body(job_id: &str) -> Vec<u8> {
1351 serde_json::to_vec(&json!({
1352 "v": AGENT_PROTOCOL_VERSION,
1353 "op": "jobCancel",
1354 "jobId": job_id,
1355 }))
1356 .unwrap_or_default()
1357}
1358
1359fn deadline_millis(deadline: Duration) -> u64 {
1362 u64::try_from(deadline.as_millis()).unwrap_or(u64::MAX)
1363}
1364
1365fn confirm_empty_ok(operation: &str, body: &[u8]) -> Result<()> {
1370 if body.iter().all(|byte| byte.is_ascii_whitespace()) {
1371 return Ok(());
1372 }
1373 Err(AlienError::new(ErrorData::SandboxCommandFailed {
1374 failure: "agentRefused".to_string(),
1375 reason: format!("{operation} was refused: {}", truncated(body)),
1376 }))
1377}
1378
1379fn session_segment(name: &str) -> Option<&str> {
1381 let segment = name.rsplit('/').next()?;
1382 is_addressable_id(segment).then_some(segment)
1383}
1384
1385fn is_addressable_id(id: &str) -> bool {
1386 !id.is_empty()
1387 && id.len() <= MAX_SESSION_ID
1388 && id
1389 .chars()
1390 .all(|c| c.is_ascii_alphanumeric() || c == '-' || c == '_')
1391}
1392
1393fn generation_from_boot_id(boot_id: &str) -> u64 {
1399 const FNV_OFFSET_BASIS: u64 = 0xcbf2_9ce4_8422_2325;
1400 const FNV_PRIME: u64 = 0x0000_0100_0000_01b3;
1401 let mut hash = FNV_OFFSET_BASIS;
1402 for byte in boot_id.as_bytes() {
1403 hash ^= u64::from(*byte);
1404 hash = hash.wrapping_mul(FNV_PRIME);
1405 }
1406 hash | 1
1407}
1408
1409fn session_state(operation: &str, state: Option<&str>) -> Result<SandboxSessionState> {
1412 match state {
1413 Some("STATE_RUNNING") => Ok(SandboxSessionState::Running),
1414 Some("STATE_CREATING" | "STATE_PENDING" | "STATE_RESUMING") => {
1415 Ok(SandboxSessionState::Starting)
1416 }
1417 Some("STATE_PAUSED" | "STATE_PAUSING" | "STATE_SUSPENDED") => {
1418 Ok(SandboxSessionState::Suspended)
1419 }
1420 Some("STATE_STOPPED" | "STATE_FAILED" | "STATE_DELETING" | "STATE_DELETED") => {
1421 Ok(SandboxSessionState::Terminated)
1422 }
1423 other => Err(AlienError::new(ErrorData::UnexpectedResponseFormat {
1424 provider: "gcp-agent-platform".to_string(),
1425 binding_name: operation.to_string(),
1426 field: "state".to_string(),
1427 response_json: other
1428 .map_or_else(|| "absent".to_string(), |state| format!("\"{state}\"")),
1429 })),
1430 }
1431}
1432
1433fn finish_operation(operation: &str, name: &str, op: Operation) -> Result<serde_json::Value> {
1435 match op.result {
1436 Some(OperationResult::Response { response }) => Ok(response),
1437 Some(OperationResult::Error { error }) => {
1438 Err(AlienError::new(ErrorData::SandboxCommandFailed {
1439 failure: "operationFailed".to_string(),
1440 reason: format!(
1441 "{operation}: operation '{name}' failed (grpc {}): {}",
1442 error.code, error.message
1443 ),
1444 }))
1445 }
1446 None => Err(AlienError::new(ErrorData::UnexpectedResponseFormat {
1447 provider: "gcp-agent-platform".to_string(),
1448 binding_name: operation.to_string(),
1449 field: "response".to_string(),
1450 response_json: format!("operation '{name}' reported done without a result"),
1451 })),
1452 }
1453}
1454
1455fn is_not_found(error: &AlienError<AgentPlatformErrorData>) -> bool {
1461 const NOT_FOUND: &str = "REMOTE_RESOURCE_NOT_FOUND";
1462 if error.code == NOT_FOUND {
1463 return true;
1464 }
1465 let mut node = error.source.as_deref();
1466 while let Some(current) = node {
1467 if current.code == NOT_FOUND {
1468 return true;
1469 }
1470 node = current.source.as_deref();
1471 }
1472 false
1473}
1474
1475fn truncated(body: &[u8]) -> String {
1477 const LIMIT: usize = 200;
1478 let text = String::from_utf8_lossy(body);
1479 let text = text.trim();
1480 if text.len() <= LIMIT {
1481 return text.to_string();
1482 }
1483 let end = (0..=LIMIT)
1484 .rev()
1485 .find(|at| text.is_char_boundary(*at))
1486 .unwrap_or(0);
1487 format!("{}…", &text[..end])
1488}
1489
1490#[cfg(test)]
1491#[path = "gcp_agent_platform_tests.rs"]
1492mod tests;