Skip to main content

alien_bindings/providers/sandbox/
gcp_agent_platform.rs

1//! GCP Agent Platform sandbox provider.
2//!
3//! Sandboxes are `sandboxEnvironments` created under a durable reasoning engine and reached from
4//! outside the guest through the `:execute` proxy, which forwards one request to the agent's
5//! `POST /` envelope and returns its body verbatim. So every command, file operation and health
6//! check is one envelope over that proxy, and the lifecycle verbs are long-running operations
7//! polled to completion.
8
9use 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
34/// The envelope protocol version this provider speaks. It matches the agent's `PROTOCOL_VERSION`;
35/// a peer that answers a different one is refused rather than guessed at.
36const AGENT_PROTOCOL_VERSION: u32 = 1;
37
38/// The proxy holds one `:execute` request open for roughly this long, so a command whose deadline
39/// is within it runs synchronously and anything longer is detached as a job and polled. Set below
40/// the measured ceiling, because a command that overruns a synchronous execute is lost, where an
41/// overrun job is still reachable by a later poll.
42const MAX_SYNCHRONOUS_TIMEOUT: Duration = Duration::from_secs(30);
43
44/// Longest sandbox id this provider will place in a proxy URL. A bound on what is handed back to a
45/// caller, not on what the API mints — the names seen are far shorter.
46const MAX_SANDBOX_ID: usize = 63;
47
48/// How long a created sandbox has to reach `STATE_RUNNING`, and how often that is checked.
49const SANDBOX_READY_ATTEMPTS: u32 = 150;
50const SANDBOX_READY_INTERVAL: Duration = Duration::from_secs(2);
51
52/// How long `create` keeps probing a RUNNING sandbox's agent, and how often. The proxy answers
53/// `Bad Gateway` for a few seconds after RUNNING while the container starts listening.
54#[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// A unit test's mock answers at once, so the budget only decides how long it takes to say no.
59#[cfg(test)]
60const AGENT_READY_TIMEOUT: Duration = Duration::from_millis(50);
61#[cfg(test)]
62const AGENT_READY_POLL: Duration = Duration::from_millis(5);
63
64/// How long a lifecycle operation (`create`, `:snapshot`) is polled before it
65/// is reported incomplete rather than waited on forever.
66const OPERATION_POLL_ATTEMPTS: u32 = 150;
67const OPERATION_POLL_INTERVAL: Duration = Duration::from_secs(2);
68
69/// How long `terminate` polls the sandbox to `not-found`, turning an accepted delete into a
70/// confirmed one.
71const TERMINATE_POLL_ATTEMPTS: u32 = 30;
72const TERMINATE_POLL_INTERVAL: Duration = Duration::from_secs(2);
73
74/// How often a detached job is polled for new output.
75const JOB_POLL_INTERVAL: Duration = Duration::from_secs(1);
76
77/// The grace a job's poll loop allows past the command's own deadline before it cancels the job:
78/// the agent kills the command at the deadline and the next poll reports it, and this covers the
79/// round trips to observe that.
80const 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
91/// The generation of a sandbox whose live container identity was not established: a state with no
92/// reachable agent, or a bulk `list` that does not probe each sandbox. Never a value
93/// `generation_from_boot_id` returns, so a real identity is always distinguishable from an
94/// unprobed one.
95const NO_GENERATION: u64 = 0;
96
97/// A single health probe is bounded to this, because the client sets no per-request timeout and an
98/// agent that accepts the connection but never answers would otherwise hang `get()` and `create()`
99/// forever. Set above the proxy's ~30s synchronous window (see `MAX_SYNCHRONOUS_TIMEOUT`) rather
100/// than tight to the round trip: too tight reports a healthy sandbox unreachable, and
101/// `get_or_create` then provisions a fresh sandbox and loses the caller's filesystem — the failure
102/// this task exists to prevent — where too loose only delays an already-broken sandbox.
103const AGENT_PROBE_BUDGET: Duration = Duration::from_secs(60);
104
105/// Maps a declared egress mode onto the template's `egressControlConfig`, or refuses one the API
106/// cannot express.
107///
108/// `internetAccess` is a single boolean, so `AllowDomains` has no representation and is refused
109/// rather than approximated into `allow` (which would open more than was asked) or `deny` (which
110/// would close a caller out of hosts it named). Not called by the runtime verbs — the template is
111/// pre-created — but this is the mapping the template controller uses, kept beside the provider so
112/// the two agree on what a mode means. `sandbox_label` names the offending sandbox in the refusal.
113pub 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/// A Sandbox backed by the Vertex AI Agent Platform.
135#[derive(Debug)]
136pub struct GcpAgentPlatformSandbox {
137    client: Arc<dyn AgentPlatformApi>,
138    /// Bare reasoning-engine id the client interpolates into its paths. The binding may carry a
139    /// full resource name, so it is reduced to its last segment once, here.
140    engine: String,
141    /// Template every sandbox is cut from, as a resource name the create body carries unchanged.
142    template: String,
143    /// Sandbox lifetime in seconds, from the declaration; absent takes the service default.
144    max_lifetime_seconds: Option<u32>,
145}
146
147impl GcpAgentPlatformSandbox {
148    /// The `ttl` a create is sent with: what the caller asked for, never above what the
149    /// declaration allows.
150    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    /// Builds a provider bound to one engine and template.
161    ///
162    /// The engine is normalised to its last path segment because the client builds the full
163    /// resource path itself; passing the whole name would double it and address nothing.
164    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    /// The engine id sent to the client. Exists so a test can prove the binding's full resource
180    /// name was reduced to a bare segment — a doubled path is invisible against the mock otherwise.
181    #[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    /// Refused rather than forwarded: Agent Platform's `:resume` can return a fresh container while
194    /// reporting success, so nothing the sandbox held is guaranteed to survive a pause.
195    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    /// A sandbox id that stays a single path segment.
203    ///
204    /// The id is interpolated into the proxy URL, so one carrying `/`, `..`, `?` or `#` would
205    /// address a different sandbox — a resource the same engine grant can reach. The API mints
206    /// these; this bounds the ones a caller hands back.
207    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    /// Reads a sandbox, or `None` when it is gone, without judging it.
222    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    /// Polls a lifecycle operation to completion, returning its response payload.
238    ///
239    /// Bounded rather than open-ended: a caller waiting forever is its own outage, and the
240    /// operation name is carried so an incomplete one can be resumed rather than lost.
241    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    /// Sends one envelope through the `:execute` proxy and returns the agent's body verbatim.
281    ///
282    /// A client error is a transport failure — the proxy could not deliver or the API refused. A
283    /// body the op's parser cannot read is the agent's own reason, handled by each verb. A
284    /// not-found is reported as a gone sandbox so a caller does not read it as a live one.
285    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        // The client does not tell a delivered-but-failed call apart from an undelivered one, so
320        // a `:execute` carrying a command leaves its outcome unestablished. The cause stays on the
321        // chain rather than in `reason`, keeping a redacted request body out of an externally
322        // visible message.
323        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    /// Confirms the agent answers and speaks the protocol, and returns the sandbox's generation.
336    ///
337    /// A sandbox can report `STATE_RUNNING` while every `:execute` fails, so a state read is not a
338    /// health check; the agent has to answer for the sandbox to be usable. The reply carries the
339    /// container boot id, from which the generation is derived so a caller can detect a container
340    /// that was replaced under a stable sandbox name.
341    async fn probe_agent(&self, operation: &str, sandbox_id: &str) -> Result<u64> {
342        // Mapped to unreachable whatever the failure — a refused delivery, a probe that outran its
343        // budget, an unparseable body, a protocol mismatch — because a health probe is idempotent
344        // and the caller acts on the same thing each way: the agent cannot be reached, so
345        // `get_or_create` provisions a fresh one rather than destroying a sandbox it did not create.
346        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        // An agent that answers without a boot id cannot be told apart from a replaced container,
398        // so the sandbox is refused rather than reconnected to a possibly-blank one.
399        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    /// Deletes a sandbox the caller will never receive, keeping the reason it is discarded.
410    ///
411    /// Every failure after the sandbox exists reaches here, so `create` has one delete rather than
412    /// one beside each `?`. The delete's own failure names the leak without replacing the finding
413    /// that caused it. A not-found delete is already success in the client.
414    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    /// Waits for a created sandbox to reach `STATE_RUNNING`, confirms its agent answers, and returns
437    /// the sandbox's generation.
438    ///
439    /// The running record is judged, not the create accept: a sandbox still coming up need not be
440    /// addressable yet, and reading that as a failure would delete every one that answered early.
441    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                // Waited on rather than woken: a fresh sandbox has no idle-suspend policy to pause
460                // it before its first command — the binding carries no such field — so a suspended
461                // reading here is a transient step on the way up, not a resting state to resume.
462                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    /// Probes a RUNNING sandbox's agent until it answers or [`AGENT_READY_TIMEOUT`] passes, returning
476    /// the generation or the timeout over the last finished probe's cause. Every failure is retried,
477    /// a 404 included: the proxy answers a not-yet-listening container with a relayed 404.
478    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            // Bounded by the deadline too, since one probe's own budget is as long as the wait.
483            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    /// Runs a command inside the proxy's synchronous window, streaming the buffered NDJSON body.
509    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    /// Runs a command as a detached job whose output is polled for until it ends.
521    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    /// Refuses a command the agent would refuse anyway, before a call is spent on it.
545    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        // Refused rather than defaulted, and refused where it floors to zero milliseconds too: the
554        // agent rejects a `timeoutMs` of 0, and a defaulted timeout is a hang waiting for a slow
555        // day in a sandbox running code the caller does not control.
556        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    /// The platform's row, unnarrowed. `sandboxLifetime` stays true even with no declared ttl:
575    /// the API always sets `expireTime` on output, so an undeclared sandbox still carries a
576    /// deadline the platform enforces.
577    fn capabilities(&self) -> SandboxCapabilities {
578        SandboxCapabilities::gcp_agent_platform()
579    }
580
581    async fn create(&self, request: CreateSandboxRequest) -> Result<SandboxInstance> {
582        // A sandbox takes no environment of its own: `SandboxCreateRequest` has no env field,
583        // so silently dropping one would run the caller's code without the variables it asked for.
584        // They travel per command through `run_command` instead.
585        // `OperationNotSupported`, not `InvalidInput`: the value is fine, the backend has nowhere
586        // to put it. AWS answers the identical condition the same way, and a portable caller
587        // branching on the code must not get two answers for one situation.
588        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        // Same reason as `env` above: nowhere to carry a tenant key, so accepting one would
598        // silently merge tenants into one sandbox.
599        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        // The caller's requested id is not authoritative — the API allocates the name, and the
642        // last segment is the id every later verb addresses it by. One this client cannot send is
643        // one nothing can reach or reap, so an unreadable name is reported without a delete it
644        // cannot target.
645        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        // Past here a sandbox exists the caller has no id for, so every failure deletes it.
656        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        // Only a running sandbox carries a reachable agent, and a state read is not health: a
674        // running record whose agent does not answer is not reported as usable. A non-running
675        // sandbox has no live container to identify, so it carries no generation.
676        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            // A running, reachable sandbox is handed back; anything else is served by a fresh
692            // sandbox rather than by destroying one this call did not create, which may be
693            // another revision's.
694            match self.get(id).await {
695                Ok(Some(sandbox)) if sandbox.state == SandboxState::Running => {
696                    return Ok(ResolvedSandbox::found(sandbox))
697                }
698                // Still coming up, or already being woken by someone else. Waited for rather than
699                // replaced: the sandbox keeps starting either way, so a second one beside it would
700                // run with nobody holding its id. A slow data plane is answered with the failure.
701                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        // A sandbox this provider cannot fully read — an unaddressable name or an unrecognised
732        // state — is left out rather than surfaced as a handle to nothing or failing the whole
733        // enumeration; one odd sandbox must not hide every other from an orphan sweep. Both halves
734        // are skipped for the same reason, so leniency is consistent across the record.
735        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                // A bulk list does not probe each agent, so it reports no generation; a caller that
741                // needs one reads the single sandbox through `get`.
742                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        // The synchronous window is the proxy's, not the command's: a command that outlives one
760        // `:execute` is detached as a job so a later poll can still reach its output.
761        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        // The `:execute` succeeded, so the job was accepted and is running; only its id could not
776        // be read. Nothing can poll or cancel it after this, and the command's own deadline is
777        // what bounds it — so the outcome is unestablished rather than a reply that failed to read.
778        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        // One request per path, stopping at the first failure — the partial application every
886        // backend performs, so a caller sees one contract rather than several. The agent's field
887        // is `contentsBase64`; `contents` is dropped silently.
888        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        // A generated display name, because the API takes one and the caller does not supply it.
926        // The trait has no restore verb, so the returned name is not yet consumable through it —
927        // restore is `create` from a snapshot, which this backend can do but the trait cannot ask.
928        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        // Accepted, not completed: the client returns before the sandbox is gone. Returning here
964        // would report containment while the code may still run, which is the whole point of
965        // terminate — so the delete is confirmed by polling to not-found.
966        // A sandbox that is already gone is the state terminate exists to reach, so not-found is
967        // success. Narrowed to exactly that: mapping any failure to `Ok` would report containment
968        // for a sandbox another deployment owns and this one was refused.
969        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                // A read that fails is not a sandbox that is gone, and one throttled response must
983                // not end the poll: the attempt budget decides.
984                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
1003/// One step of a detached job's poll loop, yielding output frames as they arrive and a terminal
1004/// item once the job ends.
1005async 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            // The deadline is this client's decision, so it only names an outcome once the cancel
1016            // that makes it true has landed. A cancel that fails leaves the job running, and the
1017            // caller has to be told that rather than that the command was stopped.
1018            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            // Only a landed cancel stops the job; the other two arms report an outcome nobody
1046            // established, so the drop is left armed to try again.
1047            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            // A standalone poll is repeatable, but this one watches a running command, and giving
1062            // up on it leaves that command's outcome unestablished.
1063            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            // A seq gap is truncated output, not a frame still to come, so the cursor takes the
1077            // highest seq seen and the loop never waits for a "missing" one; `max` rather than the
1078            // last frame's seq so an out-of-order frame cannot walk the cursor backwards.
1079            state.since_seq = state.since_seq.max(frame.seq());
1080            let output = frame.into_output();
1081            // As in the synchronous path: a frame that will not convert ends the poll rather than
1082            // being queued ahead of a terminal result that would contradict it.
1083            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            // The terminal outcome is the envelope's, not a frame's: a clean exit carries a code,
1093            // and a deadline, spawn failure or cancel carries an error object with no code.
1094            let terminal = match poll.error {
1095                Some(error) => Err(AlienError::new(ErrorData::SandboxCommandFailed {
1096                    failure: error.code,
1097                    reason: error.message,
1098                })),
1099                // A job that finished without an exit code never established its outcome; an
1100                // invented code is indistinguishable from one the command really exited with.
1101                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
1124/// One `jobPoll` against a sandbox, classified as a standalone poll: nothing about the job
1125/// changes, so a call that fails is worth repeating. `run_command`'s loop re-contexts it.
1126async 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
1148/// A poll or cancel that did not complete. Both leave the job exactly as it was, so unlike a
1149/// command they carry the retry signal; a sandbox that is gone is an answer rather than a failure.
1150fn 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
1166/// Whether a `jobCancel` reply is the cancel landing.
1167///
1168/// A reply arriving is not the cancel succeeding: a cancelled job answers with the empty object and
1169/// a refusal — `JobNotFound`, say — with the agent's error text, both through a successful
1170/// `:execute`. A success body that grew a field then reads as unknown, never a refusal as landed.
1171fn 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/// The id a started job answers to, as the agent's `jobStart` returns it.
1177#[derive(Deserialize)]
1178#[serde(rename_all = "camelCase")]
1179struct JobStartReply {
1180    job_id: String,
1181}
1182
1183/// The bookkeeping a detached job's poll loop carries between steps.
1184struct 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    /// Set only where the job itself ended, so a stream that ends any other way still cancels.
1193    stopped: bool,
1194    deadline_at: tokio::time::Instant,
1195}
1196
1197/// Kills the command a dropped stream stops reading.
1198///
1199/// Where a command's frames ride a transport the caller can close, that close is the kill. A
1200/// detached job has no such transport, so the cancel is sent here or the job runs to its full
1201/// timeout, billing, with nobody able to reach it.
1202impl 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            // The reader is gone, so a failure has nobody to be returned to; what it leaves
1222            // running is named instead.
1223            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/// A job's output so far, and how it ended once it has. Mirrors the agent's `jobPoll` reply.
1254#[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/// A frame as the agent writes it, shared by the synchronous NDJSON body and the job frames.
1276#[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            // An error frame is the command's outcome, so it surfaces as an error rather than a
1322            // stream that simply stopped.
1323            Self::Error { code, message } => {
1324                Err(AlienError::new(ErrorData::SandboxCommandFailed {
1325                    failure: code,
1326                    reason: message,
1327                }))
1328            }
1329        }
1330    }
1331}
1332
1333/// A frame that arrived is proof the command ran, so a payload that will not decode leaves the
1334/// outcome unestablished rather than merely malformed.
1335fn 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
1350/// Turns the agent's buffered NDJSON body into output frames.
1351///
1352/// A body that is not frames at all is the agent's error, reported as a refusal. A body that ends
1353/// without a terminal frame is a transport failure, not a command that finished: the command had
1354/// started, so the trailing item says the outcome is unknown rather than letting a truncated
1355/// stream read as success.
1356fn 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                // A frame that will not convert ends the body: letting a later exit follow would
1371                // answer the question this item just reported as unanswerable. `saw_terminal`
1372                // stops the trailing item below from saying the same thing twice.
1373                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 already arrived, so the command ran and this leaves its end unknown.
1388                // `saw_terminal` stops the trailing item below from saying the same thing twice.
1389                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
1420/// The envelope for `exec` or `jobStart`. `timeoutMs` is the field the agent reads; both ops take
1421/// the identical body.
1422fn 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
1452/// Milliseconds, saturated: a timeout long enough to overflow `u64` ms is not one anyone meant,
1453/// and wrapping it would turn "effectively forever" into "immediately".
1454fn timeout_millis(timeout: Duration) -> u64 {
1455    u64::try_from(timeout.as_millis()).unwrap_or(u64::MAX)
1456}
1457
1458/// Reads a `writeFile` reply, which succeeds with an empty body.
1459///
1460/// A non-empty body from the op is the agent's error text, not a success shape, so it is
1461/// surfaced as a refusal rather than ignored.
1462fn 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
1472/// The last path segment, if it is a usable id. Used for both minted names and listed ones.
1473fn 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
1486/// Maps a container boot id to a numeric generation deterministically.
1487///
1488/// A caller may compare generations across processes, so this is an explicit FNV-1a rather than a
1489/// `Hash` impl — the same boot id must yield the same number in any build, and std's hashers
1490/// promise no cross-release stability. `| 1` keeps the result clear of `NO_GENERATION`.
1491fn 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
1502/// The API's runtime states, in ours. An unrecognised one is an error rather than a default,
1503/// because every default here is a lie a caller acts on.
1504fn 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
1522/// Turns a completed operation into its response payload, or the error it reported.
1523fn 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
1544/// The agent's own refusal, when `:execute` relayed one. The proxy forwards the agent's status, so
1545/// its `PATH_NOT_FOUND` is a 404 like a missing sandbox; only a relayed body carries the agent's
1546/// error code after `Error Details:`. A relayed 5xx is not an answer: the command may have run.
1547fn 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
1557/// Whether a client error means the sandbox is already gone.
1558///
1559/// The client wraps a 404 as `RequestFailed` and leaves the `RemoteResourceNotFound` on the
1560/// source chain, so the classification is read by walking that chain rather than off the outer
1561/// variant — a path or trace id mentioning 404 in a message never reaches this.
1562fn 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
1577/// A body short enough to sit in an error message without carrying a whole response into it.
1578fn 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;