Skip to main content

alien_bindings/providers/sandbox/
kubernetes.rs

1//! Kubernetes sandbox provider: a pod under a sandboxed runtime class, reached over the agent
2//! protocol.
3//!
4//! The application never holds a cluster credential. It asks the operator's broker for a
5//! session, and gets back a pod address plus a capability scoped to that session. Claiming a pod
6//! is a `PATCH` on pods, which does not belong in the binding: `pods/exec`
7//! would reach every pod in the namespace.
8//!
9//! It authenticates to the broker with the ServiceAccount token Kubernetes already mounted in
10//! its pod. Nothing of Alien's is created, rotated or torn down for this.
11
12use std::collections::BTreeMap;
13use std::sync::Mutex;
14
15use async_trait::async_trait;
16use futures::stream::BoxStream;
17use serde::{Deserialize, Serialize};
18
19use crate::error::{ErrorData, Result};
20use crate::providers::sandbox::agent_protocol::{self, AgentTransport};
21use crate::traits::{
22    Binding, CommandOutput, CreateSessionRequest, JobPoll, JobStart, PreviewCapability,
23    RunCommandRequest, Sandbox, SandboxSession, SandboxSessionState,
24};
25use alien_core::bindings::KubernetesSandboxBinding;
26use alien_core::{Platform, SandboxCapabilities};
27use alien_error::{AlienError, Context, IntoAlienError};
28
29/// What the broker hands back for a claimed session.
30#[derive(Debug, Clone, Deserialize)]
31#[serde(rename_all = "camelCase")]
32struct ClaimResponse {
33    session_id: String,
34    endpoint: String,
35    capability: String,
36    expires_at: i64,
37}
38
39#[derive(Debug, Serialize)]
40#[serde(rename_all = "camelCase")]
41struct ClaimRequest<'a> {
42    sandbox_id: &'a str,
43    session_id: &'a str,
44}
45
46/// A Sandbox backed by pods under a sandboxed runtime class.
47#[derive(Debug)]
48pub struct KubernetesSandbox {
49    sandbox_id: String,
50    broker_url: String,
51    token_path: String,
52    binding_name: String,
53    client: reqwest::Client,
54    /// Claims this process has made, so a later call can address the session it already has.
55    ///
56    /// The capability is short-lived and the endpoint is a pod IP, so this is a cache of live
57    /// sessions rather than durable state. A session this process did not claim is not
58    /// reachable, which is what `reconnect` means here.
59    claims: Mutex<BTreeMap<String, ClaimResponse>>,
60}
61
62impl KubernetesSandbox {
63    /// Builds a provider from its binding.
64    pub fn new(
65        binding_name: &str,
66        binding: &KubernetesSandboxBinding,
67        sandbox_id: &str,
68    ) -> Result<Self> {
69        let value = |field: &'static str, value: alien_core::bindings::BindingValue<String>| {
70            value.into_value(binding_name, field).map_err(|error| {
71                AlienError::new(ErrorData::BindingConfigInvalid {
72                    binding_name: binding_name.to_string(),
73                    env_var: alien_core::bindings::binding_env_var_name(binding_name),
74                    reason: error.to_string(),
75                })
76            })
77        };
78
79        Ok(Self {
80            sandbox_id: sandbox_id.to_string(),
81            broker_url: value("brokerUrl", binding.broker_url.clone())?
82                .trim_end_matches('/')
83                .to_string(),
84            token_path: value("tokenPath", binding.token_path.clone())?,
85            binding_name: binding_name.to_string(),
86            client: reqwest::Client::new(),
87            claims: Mutex::new(BTreeMap::new()),
88        })
89    }
90
91    /// Reads the pod's ServiceAccount token.
92    ///
93    /// Read per call rather than cached: Kubernetes rotates projected tokens in place, and a
94    /// cached copy becomes a token the apiserver refuses at the least convenient moment.
95    async fn identity_token(&self) -> Result<String> {
96        tokio::fs::read_to_string(&self.token_path)
97            .await
98            .into_alien_error()
99            .context(ErrorData::BindingConfigInvalid {
100                binding_name: self.binding_name.clone(),
101                env_var: alien_core::bindings::binding_env_var_name(&self.binding_name),
102                reason: format!(
103                    "could not read the ServiceAccount token at '{}'",
104                    self.token_path
105                ),
106            })
107    }
108
109    fn claimed(&self, session_id: &str) -> Option<ClaimResponse> {
110        self.claims
111            .lock()
112            .expect("no panic holds this lock")
113            .get(session_id)
114            .cloned()
115    }
116
117    fn failed(&self, operation: &str, reason: &str) -> AlienError<ErrorData> {
118        AlienError::new(ErrorData::OperationNotSupported {
119            operation: operation.to_string(),
120            reason: reason.to_string(),
121        })
122    }
123}
124
125#[async_trait]
126impl AgentTransport for KubernetesSandbox {
127    async fn request(
128        &self,
129        session_id: &str,
130        method: reqwest::Method,
131        path: &str,
132    ) -> Result<reqwest::RequestBuilder> {
133        let claim = self.claimed(session_id).ok_or_else(|| {
134            self.failed(
135                "sandbox.agent",
136                &format!(
137                    "session '{session_id}' was not claimed by this process; a pod IP and a \
138                     capability are only reachable by the caller that claimed them"
139                ),
140            )
141        })?;
142
143        if claim.expires_at <= chrono::Utc::now().timestamp() {
144            return Err(self.failed(
145                "sandbox.agent",
146                &format!(
147                    "the capability for session '{session_id}' expired; the agent would refuse \
148                     this with a 401 that reads like a broken sandbox"
149                ),
150            ));
151        }
152
153        Ok(self
154            .client
155            .request(method, format!("{}{path}", claim.endpoint))
156            .bearer_auth(claim.capability))
157    }
158
159    fn provider(&self) -> &'static str {
160        "kubernetes-sandbox"
161    }
162}
163
164impl Binding for KubernetesSandbox {}
165
166#[async_trait]
167impl Sandbox for KubernetesSandbox {
168    fn as_any(&self) -> &dyn std::any::Any {
169        self
170    }
171
172    fn capabilities(&self) -> SandboxCapabilities {
173        SandboxCapabilities::for_platform(Platform::Kubernetes)
174            .expect("Kubernetes has a sandbox backend")
175    }
176
177    /// Claims a warm pod through the broker.
178    async fn create(&self, request: CreateSessionRequest) -> Result<SandboxSession> {
179        let session_id = request
180            .session_id
181            .unwrap_or_else(|| uuid::Uuid::new_v4().simple().to_string());
182
183        let response = self
184            .client
185            .post(format!("{}/v1/sandbox/sessions", self.broker_url))
186            .bearer_auth(self.identity_token().await?)
187            .json(&ClaimRequest {
188                sandbox_id: &self.sandbox_id,
189                session_id: &session_id,
190            })
191            .send()
192            .await
193            .into_alien_error()
194            .context(ErrorData::OperationNotSupported {
195                operation: "sandbox.create".to_string(),
196                reason: "the sandbox broker is unreachable".to_string(),
197            })?;
198
199        if !response.status().is_success() {
200            let status = response.status();
201            let body = response.text().await.unwrap_or_default();
202            // 503 is the pool being empty, which refills on the controller's next health tick.
203            // Saying so is the difference between a caller retrying and a caller giving up.
204            return Err(self.failed(
205                "sandbox.create",
206                &format!("the sandbox broker returned {status}: {body}"),
207            ));
208        }
209
210        let claim: ClaimResponse = response.json().await.into_alien_error().context(
211            ErrorData::UnexpectedResponseFormat {
212                provider: "kubernetes-sandbox".to_string(),
213                binding_name: "sandbox.create".to_string(),
214                field: "body".to_string(),
215                response_json: "the broker returned a body this provider cannot parse".to_string(),
216            },
217        )?;
218
219        self.claims
220            .lock()
221            .expect("no panic holds this lock")
222            .insert(claim.session_id.clone(), claim.clone());
223
224        Ok(SandboxSession {
225            session_id: claim.session_id,
226            state: SandboxSessionState::Running,
227            // A released pod is deleted rather than fenced, so a session never outlives its own
228            // generation.
229            generation: 1,
230        })
231    }
232
233    /// Only sessions this process claimed are addressable.
234    ///
235    /// A capability is minted to the caller that claimed the pod, so another process holding the
236    /// same session id has nothing to reach it with. Returning `None` rather than erroring: the
237    /// session may well exist, this caller simply cannot address it.
238    async fn get(&self, session_id: &str) -> Result<Option<SandboxSession>> {
239        Ok(self.claimed(session_id).map(|claim| SandboxSession {
240            session_id: claim.session_id,
241            state: SandboxSessionState::Running,
242            generation: 1,
243        }))
244    }
245
246    async fn get_or_create(&self, request: CreateSessionRequest) -> Result<SandboxSession> {
247        if let Some(id) = request.session_id.as_deref() {
248            if let Some(existing) = self.get(id).await? {
249                return Ok(existing);
250            }
251        }
252
253        self.create(request).await
254    }
255
256    async fn list(&self) -> Result<Vec<SandboxSession>> {
257        Ok(self
258            .claims
259            .lock()
260            .expect("no panic holds this lock")
261            .values()
262            .map(|claim| SandboxSession {
263                session_id: claim.session_id.clone(),
264                state: SandboxSessionState::Running,
265                generation: 1,
266            })
267            .collect())
268    }
269
270    async fn run_command(
271        &self,
272        session_id: &str,
273        request: RunCommandRequest,
274    ) -> Result<BoxStream<'static, Result<CommandOutput>>> {
275        agent_protocol::run_command(self, session_id, request).await
276    }
277
278    async fn start_job(&self, session_id: &str, request: RunCommandRequest) -> Result<JobStart> {
279        agent_protocol::start_job(self, session_id, request).await
280    }
281
282    async fn poll_job(
283        &self,
284        session_id: &str,
285        job_id: &str,
286        since_seq: Option<u64>,
287    ) -> Result<JobPoll> {
288        agent_protocol::poll_job(self, session_id, job_id, since_seq).await
289    }
290
291    async fn cancel_job(&self, session_id: &str, job_id: &str) -> Result<()> {
292        agent_protocol::cancel_job(self, session_id, job_id).await
293    }
294
295    async fn read_file(&self, session_id: &str, path: &str) -> Result<Vec<u8>> {
296        agent_protocol::read_file(self, session_id, path).await
297    }
298
299    async fn write_files(&self, session_id: &str, files: BTreeMap<String, Vec<u8>>) -> Result<()> {
300        agent_protocol::write_files(self, session_id, files).await
301    }
302
303    async fn mkdir(&self, session_id: &str, path: &str) -> Result<()> {
304        agent_protocol::mkdir(self, session_id, path).await
305    }
306
307    async fn preview(&self, _session_id: &str, _port: u16) -> Result<PreviewCapability> {
308        Err(self.failed(
309            "preview",
310            "preview needs a gateway that validates a session-and-port capability, and that \
311             gateway does not exist yet",
312        ))
313    }
314
315    async fn suspend(&self, _session_id: &str) -> Result<()> {
316        Err(self.failed("suspendResume", "a pod cannot be suspended and resumed"))
317    }
318
319    async fn resume(&self, _session_id: &str) -> Result<()> {
320        Err(self.failed("suspendResume", "a pod cannot be suspended and resumed"))
321    }
322
323    async fn snapshot(&self, _session_id: &str) -> Result<String> {
324        Err(self.failed("snapshot", "a pod has no snapshot primitive"))
325    }
326
327    /// Releases the session, which deletes its pod.
328    ///
329    /// Idempotent: a session this process never claimed is already in the desired end state.
330    async fn terminate(&self, session_id: &str) -> Result<()> {
331        let Some(claim) = self.claimed(session_id) else {
332            return Ok(());
333        };
334
335        let response = self
336            .client
337            .delete(format!(
338                "{}/v1/sandbox/{}/sessions/{}",
339                self.broker_url, self.sandbox_id, claim.session_id
340            ))
341            .bearer_auth(self.identity_token().await?)
342            .send()
343            .await
344            .into_alien_error()
345            .context(ErrorData::OperationNotSupported {
346                operation: "sandbox.terminate".to_string(),
347                reason: "the sandbox broker is unreachable".to_string(),
348            })?;
349
350        if !response.status().is_success() {
351            let status = response.status();
352            let body = response.text().await.unwrap_or_default();
353            return Err(self.failed(
354                "sandbox.terminate",
355                &format!("the sandbox broker returned {status}: {body}"),
356            ));
357        }
358
359        self.claims
360            .lock()
361            .expect("no panic holds this lock")
362            .remove(session_id);
363
364        Ok(())
365    }
366}