alien_bindings/providers/sandbox/
kubernetes.rs1use 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#[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#[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: Mutex<BTreeMap<String, ClaimResponse>>,
60}
61
62impl KubernetesSandbox {
63 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 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 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 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 generation: 1,
230 })
231 }
232
233 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 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}