Skip to main content

alien_bindings/providers/sandbox/
gcp_agent_platform.rs

1//! GCP Agent Platform sandbox provider.
2//!
3//! Sessions 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, 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
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_DEADLINE: Duration = Duration::from_secs(30);
43
44/// Longest session 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_SESSION_ID: usize = 63;
47
48/// How long a created sandbox has to reach `STATE_RUNNING`, and how often that is checked.
49const SESSION_READY_ATTEMPTS: u32 = 150;
50const SESSION_READY_INTERVAL: Duration = Duration::from_secs(2);
51
52/// How long a lifecycle operation (`create`, `:pause`, `:resume`, `:snapshot`) is polled before it
53/// is reported incomplete rather than waited on forever.
54const OPERATION_POLL_ATTEMPTS: u32 = 150;
55const OPERATION_POLL_INTERVAL: Duration = Duration::from_secs(2);
56
57/// How long `terminate` polls the sandbox to `not-found`, turning an accepted delete into a
58/// confirmed one.
59const TERMINATE_POLL_ATTEMPTS: u32 = 30;
60const TERMINATE_POLL_INTERVAL: Duration = Duration::from_secs(2);
61
62/// How often a detached job is polled for new output.
63const JOB_POLL_INTERVAL: Duration = Duration::from_secs(1);
64
65/// The grace a job's poll loop allows past the command's own deadline before it cancels the job:
66/// the agent kills the command at the deadline and the next poll reports it, and this covers the
67/// round trips to observe that.
68const 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
79/// The generation of a session whose live container identity was not established: a state with no
80/// reachable agent, or a bulk `list` that does not probe each session. Never a value
81/// `generation_from_boot_id` returns, so a real identity is always distinguishable from an
82/// unprobed one.
83const NO_GENERATION: u64 = 0;
84
85/// A single health probe is bounded to this, because the client sets no per-request timeout and an
86/// agent that accepts the connection but never answers would otherwise hang `get()` and `create()`
87/// forever. Set above the proxy's ~30s synchronous window (see `MAX_SYNCHRONOUS_DEADLINE`) rather
88/// than tight to the round trip: too tight reports a healthy session unreachable, and
89/// `get_or_create` then provisions a fresh sandbox and loses the caller's filesystem — the failure
90/// this task exists to prevent — where too loose only delays an already-broken session.
91const AGENT_PROBE_BUDGET: Duration = Duration::from_secs(60);
92
93/// Maps a declared egress mode onto the template's `egressControlConfig`, or refuses one the API
94/// cannot express.
95///
96/// `internetAccess` is a single boolean, so `AllowDomains` has no representation and is refused
97/// rather than approximated into `allow` (which would open more than was asked) or `deny` (which
98/// would close a caller out of hosts it named). Not called by the runtime verbs — the template is
99/// pre-created — but this is the mapping the template controller uses, kept beside the provider so
100/// the two agree on what a mode means. `sandbox_label` names the offending sandbox in the refusal.
101pub 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/// A Sandbox backed by the Vertex AI Agent Platform.
123#[derive(Debug)]
124pub struct GcpAgentPlatformSandbox {
125    client: Arc<dyn AgentPlatformApi>,
126    /// Bare reasoning-engine id the client interpolates into its paths. The binding may carry a
127    /// full resource name, so it is reduced to its last segment once, here.
128    engine: String,
129    /// Template every session is cut from, as a resource name the create body carries unchanged.
130    template: String,
131    /// Session lifetime in seconds, from the declaration; absent takes the service default.
132    session_ttl_seconds: Option<u32>,
133}
134
135impl GcpAgentPlatformSandbox {
136    /// Builds a provider bound to one engine and template.
137    ///
138    /// The engine is normalised to its last path segment because the client builds the full
139    /// resource path itself; passing the whole name would double it and address nothing.
140    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    /// The engine id sent to the client. Exists so a test can prove the binding's full resource
156    /// name was reduced to a bare segment — a doubled path is invisible against the mock otherwise.
157    #[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    /// A session id that stays a single path segment.
170    ///
171    /// The id is interpolated into the proxy URL, so one carrying `/`, `..`, `?` or `#` would
172    /// address a different sandbox — a resource the same engine grant can reach. The API mints
173    /// these; this bounds the ones a caller hands back.
174    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    /// Reads a sandbox, or `None` when it is gone, without judging it.
189    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    /// Polls a lifecycle operation to completion, returning its response payload.
205    ///
206    /// Bounded rather than open-ended: a caller waiting forever is its own outage, and the
207    /// operation name is carried so an incomplete one can be resumed rather than lost.
208    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    /// Sends one envelope through the `:execute` proxy and returns the agent's body verbatim.
248    ///
249    /// A client error is a transport failure — the proxy could not deliver or the API refused. A
250    /// body the op's parser cannot read is the agent's own reason, handled by each verb. A
251    /// not-found is reported as a gone session so a caller does not read it as a live one.
252    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        // The client does not tell a delivered-but-failed call apart from an undelivered one, so
281        // a `:execute` carrying a command leaves its outcome unestablished. The cause stays on the
282        // chain rather than in `reason`, keeping a redacted request body out of an externally
283        // visible message.
284        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    /// Confirms the agent answers and speaks the protocol, and returns the session's generation.
297    ///
298    /// A sandbox can report `STATE_RUNNING` while every `:execute` fails, so a state read is not a
299    /// health check; the agent has to answer for the session to be usable. The reply carries the
300    /// container boot id, from which the generation is derived so a caller can detect a container
301    /// that was replaced under a stable session name.
302    async fn probe_agent(&self, operation: &str, session_id: &str) -> Result<u64> {
303        // Mapped to unreachable whatever the failure — a refused delivery, a probe that outran its
304        // budget, an unparseable body, a protocol mismatch — because a health probe is idempotent
305        // and the caller acts on the same thing each way: the agent cannot be reached, so
306        // `get_or_create` provisions a fresh one rather than destroying a session it did not create.
307        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        // An agent that answers without a boot id cannot be told apart from a replaced container,
359        // so the session is refused rather than reconnected to a possibly-blank one.
360        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    /// Deletes a sandbox the caller will never receive, keeping the reason it is discarded.
371    ///
372    /// Every failure after the sandbox exists reaches here, so `create` has one delete rather than
373    /// one beside each `?`. The delete's own failure names the leak without replacing the finding
374    /// that caused it. A not-found delete is already success in the client.
375    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    /// Waits for a created sandbox to reach `STATE_RUNNING`, confirms its agent answers, and returns
398    /// the session's generation.
399    ///
400    /// The running record is judged, not the create accept: a sandbox still coming up need not be
401    /// addressable yet, and reading that as a failure would delete every one that answered early.
402    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                // Waited on rather than woken: a fresh sandbox has no idle-suspend policy to pause
423                // it before its first command — the binding carries no such field — so a suspended
424                // reading here is a transient step on the way up, not a resting state to resume.
425                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    /// Runs a command inside the proxy's synchronous window, streaming the buffered NDJSON body.
439    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    /// Runs a command as a detached job whose output is polled for until it ends.
451    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    /// Refuses a command the agent would refuse anyway, before a call is spent on it.
474    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        // Refused rather than defaulted, and refused where it floors to zero milliseconds too: the
483        // agent rejects a `deadlineMs` of 0, and a defaulted deadline is a hang waiting for a slow
484        // day in a session running code the caller does not control.
485        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    /// The platform's row, unnarrowed. `sessionLifetime` stays true even with no declared ttl:
504    /// the API always sets `expireTime` on output, so an undeclared session still carries a
505    /// deadline the platform enforces.
506    fn capabilities(&self) -> SandboxCapabilities {
507        SandboxCapabilities::gcp_agent_platform()
508    }
509
510    async fn create(&self, request: CreateSessionRequest) -> Result<SandboxSession> {
511        // A session inherits no per-session environment: `SandboxCreateRequest` has no env field,
512        // so silently dropping one would run the caller's code without the variables it asked for.
513        // They travel per command through `run_command` instead.
514        // `OperationNotSupported`, not `InvalidInput`: the value is fine, the backend has nowhere
515        // to put it. AWS answers the identical condition the same way, and a portable caller
516        // branching on the code must not get two answers for one situation.
517        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        // Same reason as `env` above: nowhere to carry a tenant key, so accepting one would
527        // silently merge tenants into one sandbox.
528        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        // The caller's requested id is not authoritative — the API allocates the name, and the
569        // last segment is the id every later verb addresses it by. One this client cannot send is
570        // one nothing can reach or reap, so an unreadable name is reported without a delete it
571        // cannot target.
572        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        // Past here a sandbox exists the caller has no id for, so every failure deletes it.
583        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        // Only a running session carries a reachable agent, and a state read is not health: a
601        // running record whose agent does not answer is not reported as usable. A non-running
602        // session has no live container to identify, so it carries no generation.
603        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            // A running, reachable session is handed back; anything else is served by a fresh
619            // session rather than by destroying one this call did not create, which may be
620            // another revision's.
621            match self.get(id).await {
622                Ok(Some(session)) if session.state == SandboxSessionState::Running => {
623                    return Ok(session)
624                }
625                // The ordinary resting state for a reconnect: a suspended session is woken and
626                // confirmed, and handed back if it comes up healthy. A wake this call made that
627                // cannot be confirmed is put back to sleep before a fresh session is provisioned —
628                // the paused one may be another revision's, and a second live sandbox beside it is
629                // a leak the caller never receives an id for.
630                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                                // The wake could not be undone: leaving it live beside a fresh
638                                // session is a leak the caller gets no id for. Fail so the woken
639                                // session stays identifiable rather than provisioning a second one.
640                                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        // A sandbox this provider cannot fully read — an unaddressable name or an unrecognised
676        // state — is left out rather than surfaced as a handle to nothing or failing the whole
677        // enumeration; one odd sandbox must not hide every other from an orphan sweep. Both halves
678        // are skipped for the same reason, so leniency is consistent across the record.
679        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                // A bulk list does not probe each agent, so it reports no generation; a caller that
685                // needs one reads the single session through `get`.
686                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        // The synchronous window is the proxy's, not the command's: a command that outlives one
704        // `:execute` is detached as a job so a later poll can still reach its output.
705        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        // The `:execute` succeeded, so the job was accepted and is running; only its id could not
720        // be read. Nothing can poll or cancel it after this, and the command's own deadline is
721        // what bounds it — so the outcome is unestablished rather than a reply that failed to read.
722        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        // One request per path, stopping at the first failure — the partial application every
830        // backend performs, so a caller sees one contract rather than several. The agent's field
831        // is `contentsBase64`; `contents` is dropped silently.
832        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        // A generated display name, because the API takes one and the caller does not supply it.
896        // The trait has no restore verb, so the returned name is not yet consumable through it —
897        // restore is `create` from a snapshot, which this backend can do but the trait cannot ask.
898        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        // Accepted, not completed: the client returns before the sandbox is gone. Returning here
934        // would report containment while the code may still run, which is the whole point of
935        // terminate — so the delete is confirmed by polling to not-found.
936        // A session that is already gone is the state terminate exists to reach, so not-found is
937        // success. Narrowed to exactly that: mapping any failure to `Ok` would report containment
938        // for a session another deployment owns and this one was refused.
939        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                // A read that fails is not a session that is gone, and one throttled response must
953                // not end the poll: the attempt budget decides.
954                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
973/// One step of a detached job's poll loop, yielding output frames as they arrive and a terminal
974/// item once the job ends.
975async 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            // The deadline is this client's decision, so it only names an outcome once the cancel
986            // that makes it true has landed. A cancel that fails leaves the job running, and the
987            // caller has to be told that rather than that the command was stopped.
988            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            // A standalone poll is repeatable, but this one watches a running command, and giving
1029            // up on it leaves that command's outcome unestablished.
1030            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            // A seq gap is truncated output, not a frame still to come, so the cursor takes the
1044            // highest seq seen and the loop never waits for a "missing" one; `max` rather than the
1045            // last frame's seq so an out-of-order frame cannot walk the cursor backwards.
1046            state.since_seq = state.since_seq.max(frame.seq());
1047            let output = frame.into_output();
1048            // As in the synchronous path: a frame that will not convert ends the poll rather than
1049            // being queued ahead of a terminal result that would contradict it.
1050            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            // The terminal outcome is the envelope's, not a frame's: a clean exit carries a code,
1060            // and a deadline, spawn failure or cancel carries an error object with no code.
1061            let terminal = match poll.error {
1062                Some(error) => Err(AlienError::new(ErrorData::SandboxCommandFailed {
1063                    failure: error.code,
1064                    reason: error.message,
1065                })),
1066                // A job that finished without an exit code never established its outcome; an
1067                // invented code is indistinguishable from one the command really exited with.
1068                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
1090/// One `jobPoll` against a session, classified as a standalone poll: nothing about the job
1091/// changes, so a call that fails is worth repeating. `run_command`'s loop re-contexts it.
1092async 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
1114/// A poll or cancel that did not complete. Both leave the job exactly as it was, so unlike a
1115/// command they carry the retry signal; a session that is gone is an answer rather than a failure.
1116fn 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
1132/// Whether a `jobCancel` reply is the cancel landing.
1133///
1134/// A reply arriving is not the cancel succeeding: the agent answers `{}` when it cancelled the job
1135/// and its own error text when it did not — `JobNotFound`, say — and both come back through a
1136/// successful `:execute`.
1137fn cancel_confirmed(body: &[u8]) -> bool {
1138    serde_json::from_slice::<serde_json::Value>(body).is_ok_and(|value| value.is_object())
1139}
1140
1141/// The id a started job answers to, as the agent's `jobStart` returns it.
1142#[derive(Deserialize)]
1143#[serde(rename_all = "camelCase")]
1144struct JobStartReply {
1145    job_id: String,
1146}
1147
1148/// The bookkeeping a detached job's poll loop carries between steps.
1149struct 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/// A job's output so far, and how it ended once it has. Mirrors the agent's `jobPoll` reply.
1161#[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/// A frame as the agent writes it, shared by the synchronous NDJSON body and the job frames.
1183#[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            // An error frame is the command's outcome, so it surfaces as an error rather than a
1229            // stream that simply stopped.
1230            Self::Error { code, message } => {
1231                Err(AlienError::new(ErrorData::SandboxCommandFailed {
1232                    failure: code,
1233                    reason: message,
1234                }))
1235            }
1236        }
1237    }
1238}
1239
1240/// A frame that arrived is proof the command ran, so a payload that will not decode leaves the
1241/// outcome unestablished rather than merely malformed.
1242fn 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
1257/// Turns the agent's buffered NDJSON body into output frames.
1258///
1259/// A body that is not frames at all is the agent's error, reported as a refusal. A body that ends
1260/// without a terminal frame is a transport failure, not a command that finished: the command had
1261/// started, so the trailing item says the outcome is unknown rather than letting a truncated
1262/// stream read as success.
1263fn 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                // A frame that will not convert ends the body: letting a later exit follow would
1278                // answer the question this item just reported as unanswerable. `saw_terminal`
1279                // stops the trailing item below from saying the same thing twice.
1280                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 already arrived, so the command ran and this leaves its end unknown.
1295                // `saw_terminal` stops the trailing item below from saying the same thing twice.
1296                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
1327/// The envelope for `exec` or `jobStart`. `deadlineMs` is the field the agent reads; both ops take
1328/// the identical body.
1329fn 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
1359/// Milliseconds, saturated: a deadline long enough to overflow `u64` ms is not one anyone meant,
1360/// and wrapping it would turn "effectively forever" into "immediately".
1361fn deadline_millis(deadline: Duration) -> u64 {
1362    u64::try_from(deadline.as_millis()).unwrap_or(u64::MAX)
1363}
1364
1365/// Reads a `writeFile`/`mkdir` reply, which succeeds with an empty body.
1366///
1367/// A non-empty body from these ops is the agent's error text, not a success shape, so it is
1368/// surfaced as a refusal rather than ignored.
1369fn 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
1379/// The last path segment, if it is a usable id. Used for both minted names and listed ones.
1380fn 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
1393/// Maps a container boot id to a numeric generation deterministically.
1394///
1395/// A caller may compare generations across processes, so this is an explicit FNV-1a rather than a
1396/// `Hash` impl — the same boot id must yield the same number in any build, and std's hashers
1397/// promise no cross-release stability. `| 1` keeps the result clear of `NO_GENERATION`.
1398fn 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
1409/// The API's runtime states, in ours. An unrecognised one is an error rather than a default,
1410/// because every default here is a lie a caller acts on.
1411fn 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
1433/// Turns a completed operation into its response payload, or the error it reported.
1434fn 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
1455/// Whether a client error means the sandbox is already gone.
1456///
1457/// The client wraps a 404 as `RequestFailed` and leaves the `RemoteResourceNotFound` on the
1458/// source chain, so the classification is read by walking that chain rather than off the outer
1459/// variant — a path or trace id mentioning 404 in a message never reaches this.
1460fn 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
1475/// A body short enough to sit in an error message without carrying a whole response into it.
1476fn 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;