1use std::collections::BTreeMap;
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::BoxStream;
17use futures::StreamExt;
18use tokio::sync::mpsc;
19
20use crate::error::{ErrorData, Result};
21use crate::traits::{
22 Binding, CommandOutput, CreateSessionRequest, PreviewCapability, RunCommandRequest, Sandbox,
23 SandboxSession, SandboxSessionState,
24};
25use alien_core::bindings::GcpSandboxBinding;
26use alien_core::sandbox_process::{self, ProcessFrame, ProcessStream, FRAME_CHANNEL_DEPTH};
27use alien_core::{Platform, SandboxCapabilities};
28use alien_error::AlienError;
29
30const OUTPUT_CAP: usize = 8 * 1024 * 1024;
32
33const CONTROL_DEADLINE: Duration = Duration::from_secs(60);
35
36#[derive(Debug)]
38pub struct GcpSandbox {
39 launcher_path: String,
40 allow_egress: bool,
41 binding_name: String,
42}
43
44impl GcpSandbox {
45 pub fn new(binding_name: &str, binding: &GcpSandboxBinding) -> Result<Self> {
47 let launcher_path = binding
48 .launcher_path
49 .clone()
50 .into_value(binding_name, "launcherPath")
51 .map_err(|error| {
52 AlienError::new(ErrorData::BindingConfigInvalid {
53 binding_name: binding_name.to_string(),
54 env_var: alien_core::bindings::binding_env_var_name(binding_name),
55 reason: error.to_string(),
56 })
57 })?;
58
59 let allow_egress = binding
60 .allow_egress
61 .clone()
62 .into_value(binding_name, "allowEgress")
63 .map_err(|error| {
64 AlienError::new(ErrorData::BindingConfigInvalid {
65 binding_name: binding_name.to_string(),
66 env_var: alien_core::bindings::binding_env_var_name(binding_name),
67 reason: error.to_string(),
68 })
69 })?;
70
71 Ok(Self {
72 launcher_path,
73 allow_egress,
74 binding_name: binding_name.to_string(),
75 })
76 }
77
78 async fn control(&self, operation: &str, arguments: &[String]) -> Result<Vec<u8>> {
83 let child = sandbox_process::spawn(&self.launcher_path, arguments)
84 .and_then(|mut command| command.spawn())
85 .map_err(|error| {
86 self.failed(operation, &format!("launcher would not start: {error}"))
87 })?;
88
89 let frames = sandbox_process::run(child, CONTROL_DEADLINE, OUTPUT_CAP).await;
90
91 let mut stdout = Vec::new();
92 let mut stderr = Vec::new();
93 for frame in &frames {
94 match frame {
95 ProcessFrame::Output {
96 stream: ProcessStream::Stdout,
97 data,
98 ..
99 } => stdout.extend_from_slice(data),
100 ProcessFrame::Output {
101 stream: ProcessStream::Stderr,
102 data,
103 ..
104 } => stderr.extend_from_slice(data),
105 _ => {}
106 }
107 }
108
109 match frames.last() {
110 Some(ProcessFrame::Exit { code: 0, .. }) => Ok(stdout),
111 Some(ProcessFrame::Exit { code, .. }) => Err(self.failed(
114 operation,
115 &format!(
116 "launcher exited with {code}: {}",
117 String::from_utf8_lossy(&stderr).trim()
118 ),
119 )),
120 Some(ProcessFrame::Failed { code, message }) => {
121 Err(self.failed(operation, &format!("{code}: {message}")))
122 }
123 _ => Err(self.failed(operation, "launcher produced no terminal frame")),
124 }
125 }
126
127 fn failed(&self, operation: &str, reason: &str) -> AlienError<ErrorData> {
128 AlienError::new(ErrorData::OperationNotSupported {
129 operation: operation.to_string(),
130 reason: format!("{reason} (binding '{}')", self.binding_name),
131 })
132 }
133
134 fn unsupported(&self, capability: &str, reason: &str) -> AlienError<ErrorData> {
135 AlienError::new(ErrorData::OperationNotSupported {
136 operation: capability.to_string(),
137 reason: reason.to_string(),
138 })
139 }
140
141 fn checked_path(&self, path: &str, operation: &str) -> Result<String> {
143 if path.is_empty() || path.split('/').any(|part| part == "..") {
144 return Err(self.failed(operation, &format!("path '{path}' traverses upward")));
145 }
146 Ok(path.to_string())
147 }
148
149 fn exec_arguments(&self, session_id: &str, command: &[String]) -> Vec<String> {
151 let mut arguments = vec!["exec".to_string(), session_id.to_string(), "--".to_string()];
152 arguments.extend(command.iter().cloned());
153 arguments
154 }
155}
156
157impl Binding for GcpSandbox {}
158
159#[async_trait]
160impl Sandbox for GcpSandbox {
161 fn as_any(&self) -> &dyn std::any::Any {
162 self
163 }
164
165 fn capabilities(&self) -> SandboxCapabilities {
166 SandboxCapabilities::for_platform(Platform::Gcp).expect("GCP has a sandbox backend")
167 }
168
169 async fn create(&self, request: CreateSessionRequest) -> Result<SandboxSession> {
174 if !request.env.is_empty() {
175 return Err(self.failed(
176 "sandbox.create",
177 "the Cloud Run sandbox launcher takes no session environment; bake it into the \
178 image or pass it in each command",
179 ));
180 }
181
182 let session_id = request
183 .session_id
184 .unwrap_or_else(|| uuid::Uuid::new_v4().simple().to_string());
185
186 let mut arguments = vec!["run".to_string(), "--id".to_string(), session_id.clone()];
187 if self.allow_egress {
188 arguments.push("--allow-egress".to_string());
189 }
190
191 self.control("sandbox.create", &arguments).await?;
192
193 Ok(SandboxSession {
194 session_id,
195 state: SandboxSessionState::Running,
196 generation: 1,
199 })
200 }
201
202 async fn get(&self, _session_id: &str) -> Result<Option<SandboxSession>> {
204 Err(self.unsupported(
205 "reconnect",
206 "a Cloud Run sandbox id is scoped to one instance, and session affinity held 2 of \
207 100 five-turn conversations",
208 ))
209 }
210
211 async fn get_or_create(&self, request: CreateSessionRequest) -> Result<SandboxSession> {
212 self.create(request).await
213 }
214
215 async fn list(&self) -> Result<Vec<SandboxSession>> {
216 Err(self.unsupported(
217 "reconnect",
218 "the launcher has no enumeration verb, and an id reaches only the instance that \
219 created it",
220 ))
221 }
222
223 async fn run_command(
224 &self,
225 session_id: &str,
226 request: RunCommandRequest,
227 ) -> Result<BoxStream<'static, Result<CommandOutput>>> {
228 if request.command.is_empty() {
229 return Err(self.failed("sandbox.runCommand", "command is empty"));
230 }
231
232 if request.deadline.is_zero() {
233 return Err(self.failed(
234 "sandbox.runCommand",
235 "a command must carry a non-zero deadline",
236 ));
237 }
238
239 if !request.env.is_empty() {
243 return Err(self.failed(
244 "sandbox.runCommand",
245 "the Cloud Run sandbox launcher takes no per-command environment; bake it into \
246 the image or pass it in the command",
247 ));
248 }
249
250 let mut arguments = self.exec_arguments(session_id, &request.command);
251 if let Some(directory) = &request.working_directory {
252 arguments.insert(2, directory.clone());
254 arguments.insert(2, "--workdir".to_string());
255 }
256
257 let child = sandbox_process::spawn(&self.launcher_path, &arguments)
258 .and_then(|mut command| command.spawn())
259 .map_err(|error| {
260 self.failed(
261 "sandbox.runCommand",
262 &format!("launcher would not start: {error}"),
263 )
264 })?;
265
266 let (sender, receiver) = mpsc::channel(FRAME_CHANNEL_DEPTH);
267 tokio::spawn(sandbox_process::stream(
268 child,
269 request.deadline,
270 OUTPUT_CAP,
271 sender,
272 ));
273
274 Ok(
277 futures::stream::unfold(receiver, |mut receiver| async move {
278 let frame = receiver.recv().await?;
279 let item = match frame {
280 ProcessFrame::Failed { code, message } => {
281 Err(AlienError::new(ErrorData::OperationNotSupported {
282 operation: "sandbox.runCommand".to_string(),
283 reason: format!("{code}: {message}"),
284 }))
285 }
286 other => Ok(CommandOutput::from(other)),
287 };
288 Some((item, receiver))
289 })
290 .boxed(),
291 )
292 }
293
294 async fn read_file(&self, session_id: &str, path: &str) -> Result<Vec<u8>> {
295 let path = self.checked_path(path, "sandbox.readFile")?;
296 let command = vec!["/bin/cat".to_string(), path];
297 self.control(
298 "sandbox.readFile",
299 &self.exec_arguments(session_id, &command),
300 )
301 .await
302 }
303
304 async fn write_files(&self, session_id: &str, files: BTreeMap<String, Vec<u8>>) -> Result<()> {
310 for (path, contents) in files {
311 let path = self.checked_path(&path, "sandbox.writeFiles")?;
312 let encoded = BASE64.encode(&contents);
313
314 let command = vec![
315 "/bin/sh".to_string(),
316 "-c".to_string(),
317 "mkdir -p \"$(dirname \"$2\")\" && printf %s \"$1\" | base64 -d > \"$2\""
320 .to_string(),
321 "sh".to_string(),
322 encoded,
323 path,
324 ];
325
326 self.control(
327 "sandbox.writeFiles",
328 &self.exec_arguments(session_id, &command),
329 )
330 .await?;
331 }
332
333 Ok(())
334 }
335
336 async fn mkdir(&self, session_id: &str, path: &str) -> Result<()> {
337 let path = self.checked_path(path, "sandbox.mkdir")?;
338 let command = vec!["/bin/mkdir".to_string(), "-p".to_string(), path];
339 self.control("sandbox.mkdir", &self.exec_arguments(session_id, &command))
340 .await?;
341 Ok(())
342 }
343
344 async fn preview(&self, _session_id: &str, _port: u16) -> Result<PreviewCapability> {
345 Err(self.unsupported(
346 "preview",
347 "a Cloud Run sandbox has no ingress of its own and no addressable endpoint",
348 ))
349 }
350
351 async fn suspend(&self, _session_id: &str) -> Result<()> {
352 Err(self.unsupported("suspendResume", "the launcher has no suspend verb"))
353 }
354
355 async fn resume(&self, _session_id: &str) -> Result<()> {
356 Err(self.unsupported("suspendResume", "the launcher has no resume verb"))
357 }
358
359 async fn snapshot(&self, _session_id: &str) -> Result<String> {
360 Err(self.unsupported(
361 "snapshot",
362 "`sandbox fork` produces another live sandbox rather than a durable artifact",
363 ))
364 }
365
366 async fn terminate(&self, session_id: &str) -> Result<()> {
367 self.control(
368 "sandbox.terminate",
369 &["delete".to_string(), session_id.to_string()],
370 )
371 .await?;
372 Ok(())
373 }
374}
375
376impl From<ProcessFrame> for CommandOutput {
377 fn from(frame: ProcessFrame) -> Self {
378 match frame {
379 ProcessFrame::Output {
380 seq,
381 stream: ProcessStream::Stdout,
382 data,
383 } => CommandOutput::Stdout { seq, data },
384 ProcessFrame::Output {
385 seq,
386 stream: ProcessStream::Stderr,
387 data,
388 } => CommandOutput::Stderr { seq, data },
389 ProcessFrame::Exit { code, truncated } => CommandOutput::Exit { code, truncated },
390 ProcessFrame::Failed { code, message } => {
393 unreachable!("a failed frame is mapped to an error: {code} {message}")
394 }
395 }
396 }
397}
398
399#[cfg(test)]
400mod tests {
401 use super::*;
402 use alien_core::bindings::BindingValue;
403
404 fn launcher(body: &str) -> (tempfile::TempDir, GcpSandbox) {
410 let directory = tempfile::tempdir().expect("temp dir");
411 let path = directory.path().join("sandbox");
412 std::fs::write(&path, format!("#!/bin/sh\n{body}\n")).expect("write launcher");
413
414 #[cfg(unix)]
415 {
416 use std::os::unix::fs::PermissionsExt;
417 std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o755))
418 .expect("make executable");
419 }
420
421 let sandbox = GcpSandbox::new(
422 "sbx",
423 &GcpSandboxBinding {
424 launcher_path: BindingValue::value(path.display().to_string()),
425 allow_egress: BindingValue::value(false),
426 },
427 )
428 .expect("binding is valid");
429
430 (directory, sandbox)
431 }
432
433 #[tokio::test]
434 async fn create_names_the_session_and_withholds_egress() {
435 let (_dir, sandbox) = launcher(r#"echo "$@""#);
436
437 let session = sandbox
438 .create(CreateSessionRequest {
439 session_id: Some("s1".to_string()),
440 tenant_key: None,
441 env: BTreeMap::new(),
442 })
443 .await
444 .expect("create succeeds");
445
446 assert_eq!(session.session_id, "s1");
447 assert_eq!(session.state, SandboxSessionState::Running);
448 }
449
450 #[tokio::test]
453 async fn egress_comes_from_the_binding_and_not_from_the_request() {
454 let (dir, _) = launcher(r#"echo "$@" > "$(dirname "$0")/argv""#);
455 let path = dir.path().join("sandbox");
456
457 for (allow, expected) in [(false, false), (true, true)] {
458 let sandbox = GcpSandbox::new(
459 "sbx",
460 &GcpSandboxBinding {
461 launcher_path: BindingValue::value(path.display().to_string()),
462 allow_egress: BindingValue::value(allow),
463 },
464 )
465 .expect("binding is valid");
466
467 sandbox
468 .create(CreateSessionRequest {
469 session_id: Some("s1".to_string()),
470 tenant_key: None,
471 env: BTreeMap::new(),
472 })
473 .await
474 .expect("create succeeds");
475
476 let argv = std::fs::read_to_string(dir.path().join("argv")).expect("argv recorded");
477 assert_eq!(
478 argv.contains("--allow-egress"),
479 expected,
480 "binding said allow_egress={allow}, argv was: {argv}"
481 );
482 }
483 }
484
485 #[tokio::test]
488 async fn a_failing_launcher_surfaces_its_stderr() {
489 let (_dir, sandbox) = launcher(r#"echo "quota exhausted" 1>&2; exit 7"#);
490
491 let error = sandbox
492 .create(CreateSessionRequest {
493 session_id: Some("s1".to_string()),
494 tenant_key: None,
495 env: BTreeMap::new(),
496 })
497 .await
498 .expect_err("a non-zero launcher exit is a failure");
499
500 let rendered = format!("{error:?}");
501 assert!(rendered.contains("quota exhausted"), "got: {rendered}");
502 assert!(
503 rendered.contains('7'),
504 "the exit code belongs in the error: {rendered}"
505 );
506 }
507
508 #[tokio::test]
509 async fn a_command_streams_output_and_a_real_exit_code() {
510 let (_dir, sandbox) = launcher(r#"echo hello; echo problem 1>&2; exit 3"#);
511
512 let frames: Vec<_> = sandbox
513 .run_command(
514 "s1",
515 RunCommandRequest {
516 command: vec!["/bin/true".to_string()],
517 working_directory: None,
518 env: BTreeMap::new(),
519 deadline: Duration::from_secs(10),
520 },
521 )
522 .await
523 .expect("the command runs")
524 .collect()
525 .await;
526
527 let decoded: String = frames
528 .iter()
529 .filter_map(|frame| match frame {
530 Ok(CommandOutput::Stdout { data, .. }) => {
531 Some(String::from_utf8_lossy(data).to_string())
532 }
533 _ => None,
534 })
535 .collect();
536 assert!(decoded.contains("hello"), "stdout was: {decoded}");
537
538 assert!(
539 frames
540 .iter()
541 .any(|frame| matches!(frame, Ok(CommandOutput::Stderr { .. }))),
542 "stderr must be framed, not dropped"
543 );
544
545 assert!(
546 matches!(frames.last(), Some(Ok(CommandOutput::Exit { code: 3, .. }))),
547 "the terminal frame must carry the real exit code: {:?}",
548 frames.last()
549 );
550 }
551
552 #[tokio::test]
557 async fn an_environment_the_launcher_cannot_carry_is_refused() {
558 let (_dir, sandbox) = launcher("exit 0");
559 let env = BTreeMap::from([("TOKEN".to_string(), "secret".to_string())]);
560
561 let on_create = sandbox
562 .create(CreateSessionRequest {
563 session_id: Some("s1".to_string()),
564 tenant_key: None,
565 env: env.clone(),
566 })
567 .await
568 .expect_err("a session environment cannot be honoured here");
569 assert_eq!(on_create.code, "OPERATION_NOT_SUPPORTED");
570
571 let Err(on_command) = sandbox
573 .run_command(
574 "s1",
575 RunCommandRequest {
576 command: vec!["true".to_string()],
577 working_directory: None,
578 env,
579 deadline: Duration::from_secs(5),
580 },
581 )
582 .await
583 else {
584 panic!("a command environment cannot be honoured here");
585 };
586 assert_eq!(on_command.code, "OPERATION_NOT_SUPPORTED");
587
588 sandbox
591 .create(CreateSessionRequest {
592 session_id: Some("s2".to_string()),
593 tenant_key: None,
594 env: BTreeMap::new(),
595 })
596 .await
597 .expect("a session with no environment is fine");
598 }
599
600 #[tokio::test]
603 async fn a_command_without_a_deadline_is_refused() {
604 let (_dir, sandbox) = launcher("exit 0");
605
606 let Err(error) = sandbox
607 .run_command(
608 "s1",
609 RunCommandRequest {
610 command: vec!["true".to_string()],
611 working_directory: None,
612 env: BTreeMap::new(),
613 deadline: Duration::ZERO,
614 },
615 )
616 .await
617 else {
618 panic!("a zero deadline is not a deadline");
619 };
620 assert_eq!(error.code, "OPERATION_NOT_SUPPORTED");
621 assert!(
622 error.to_string().contains("non-zero deadline"),
623 "the message should say what was wrong, got: {error}"
624 );
625 }
626
627 #[tokio::test]
629 async fn unsupported_capabilities_error_rather_than_pretend() {
630 let (_dir, sandbox) = launcher("exit 0");
631 let capabilities = sandbox.capabilities();
632
633 assert!(!capabilities.reconnect);
634 assert!(!capabilities.preview);
635 assert!(!capabilities.suspend_resume);
636 assert!(!capabilities.snapshot);
637
638 sandbox
639 .get("s1")
640 .await
641 .expect_err("reconnect is not offered");
642 sandbox
643 .list()
644 .await
645 .expect_err("enumeration is not offered");
646 sandbox
647 .preview("s1", 8080)
648 .await
649 .expect_err("preview is not offered");
650 sandbox
651 .suspend("s1")
652 .await
653 .expect_err("suspend is not offered");
654 sandbox
655 .resume("s1")
656 .await
657 .expect_err("resume is not offered");
658 sandbox
659 .snapshot("s1")
660 .await
661 .expect_err("snapshot is not offered");
662 }
663
664 #[tokio::test]
666 async fn a_traversing_path_is_refused_before_the_launcher_sees_it() {
667 let (_dir, sandbox) = launcher("exit 0");
668
669 sandbox
670 .read_file("s1", "../etc/passwd")
671 .await
672 .expect_err("a traversing path must be refused");
673 sandbox
674 .write_files(
675 "s1",
676 BTreeMap::from([("../etc/passwd".to_string(), b"x".to_vec())]),
677 )
678 .await
679 .expect_err("a traversing path must be refused on write too");
680 }
681}