Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
Rust
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265use crate::command::{self, Spec};use crate::exec;use crate::protocol::{self, Message, v1};use anyhow::{Context, Result};use nix::sys::signal::{Signal, kill};use nix::unistd::Pid;use pty_process::Size;use tokio::io::{AsyncReadExt, BufReader};use tokio_vsock::{VsockAddr, VsockStream};use tracing::{info, warn};
const READ_CHUNK: usize = 32 * 1024;
pub async fn run(host_cid: u32, open: v1::OpenDebugShell) { if let Err(error) = serve(host_cid, open).await { warn!(%error, "debug shell session failed"); }}
async fn serve(host_cid: u32, open: v1::OpenDebugShell) -> Result<()> { let mut conn = VsockStream::connect(VsockAddr::new(host_cid, open.vsock_port)) .await .with_context(|| format!("dial host debug vsock port {}", open.vsock_port))?; info!(port = open.vsock_port, "debug shell connected");
let rows = clamp_tty_dim(open.rows); let cols = clamp_tty_dim(open.cols);
let (mut pty_reader, mut pty_writer, mut child) = match debug_shell_spec(&open) .and_then(|spec| command::spawn_pty(spec, rows, cols).context("spawn pty shell")) { Ok(child) => child, Err(error) => { report_start_failure(&mut conn, &error).await; return Err(error); } }; let pid = child.id();
let (conn_reader, conn_writer) = tokio::io::split(conn); let mut conn_reader = BufReader::new(conn_reader); let mut conn_writer = conn_writer;
let mut buf = vec![0u8; READ_CHUNK]; let client_gone = loop { tokio::select! { read = pty_reader.read(&mut buf) => match read { Ok(0) => break false, // shell exited Ok(n) => { let msg = Message { id: "pty".to_owned(), pty_data: Some(v1::PtyData { data: buf[..n].to_vec().into() }), ..Default::default() }; if protocol::write_message(&mut conn_writer, &msg).await.is_err() { break true; } } Err(error) => { // linux returns EIO (not a clean EOF) on the master once the // slave side is fully closed, so treat that as the shell // exiting normally rather than a real read failure. if error.raw_os_error() != Some(nix::libc::EIO) { warn!(%error, "pty master read failed"); } break false; } }, incoming = protocol::read_message(&mut conn_reader) => match incoming { Ok(Some(msg)) => { if let Some(data) = msg.pty_data { use tokio::io::AsyncWriteExt; if pty_writer.write_all(&data.data).await.is_err() { break false; } } else if let Some(resize) = msg.pty_resize { let size = Size::new(clamp_tty_dim(resize.rows), clamp_tty_dim(resize.cols)); if let Err(error) = pty_writer.resize(size) { warn!(%error, "pty resize failed"); } } // anything else on the debug channel is ignored } Ok(None) => break true, // client closed the connection Err(error) => { warn!(%error, "debug channel read failed"); break true; } }, } };
// if the client disconnected first, hang up the shell's process group so we // don't leak a detached session. (pty-process calls setsid in the child, so // it leads a new session and process group => pgid == pid.) if client_gone && let Some(pid) = pid { let _ = kill(Pid::from_raw(-(pid as i32)), Signal::SIGHUP); }
let exit_code = match child.wait().await { Ok(status) => { use std::os::unix::process::ExitStatusExt; status .code() .or_else(|| status.signal().map(|signal| 128 + signal)) .unwrap_or(1) } Err(error) => { warn!(%error, "waiting on debug shell failed"); 1 } };
let exit = Message { id: "pty".to_owned(), exec_exit: Some(v1::ExecExit { exit_code, error: String::new(), timed_out: false, }), ..Default::default() }; let _ = protocol::write_message(&mut conn_writer, &exit).await; info!(exit_code, "debug shell session ended"); Ok(())}async fn report_start_failure(conn: &mut VsockStream, error: &anyhow::Error) { let detail = format!("{error:#}"); let (output, exit) = start_failure_messages(&detail); let _ = protocol::write_message(conn, &output).await; let _ = protocol::write_message(conn, &exit).await;}
fn start_failure_messages(detail: &str) -> (Message, Message) { let output = Message { id: "pty".to_owned(), pty_data: Some(v1::PtyData { data: format!("error: debug shell failed: {detail}\r\n") .into_bytes() .into(), }), ..Default::default() }; let exit = Message { id: "pty".to_owned(), exec_exit: Some(v1::ExecExit { exit_code: 1, error: detail.to_owned(), timed_out: false, }), ..Default::default() }; (output, exit)}
fn debug_shell_spec(open: &v1::OpenDebugShell) -> Result<Spec> { let failed_step = open .failed_step .as_ref() .context("debug shell request is missing failed step context")?; let user = exec::resolve_user(&failed_step.user).map_err(anyhow::Error::msg)?; Ok(debug_shell_spec_for_user(open, failed_step, &user))}
fn debug_shell_spec_for_user( open: &v1::OpenDebugShell, failed_step: &v1::ExecStart, user: &exec::ResolvedUser,) -> Spec { let shell = user.shell().to_string_lossy().into_owned(); let command = format!("{}exec \"$0\" -l", open.shell_prelude);
let mut request = failed_step.clone(); request.argv = vec![shell.clone(), "-lc".to_owned(), command, shell]; request.timeout_seconds = 0; let mut spec = exec::request_spec_for_user(&request, user); let term = if open.term.is_empty() { "xterm-256color" } else { &open.term }; spec.env.push(("TERM".into(), term.into())); spec}
fn clamp_tty_dim(value: u32) -> u16 { value.clamp(1, u16::MAX as u32) as u16}
#[cfg(test)]mod tests { use std::collections::HashMap; use std::ffi::{OsStr, OsString}; use std::path::Path;
use super::*;
#[test] fn debug_shell_start_failure_is_sent_to_client() { let detail = "debug shell request is missing failed step context"; let (output, exit) = start_failure_messages(detail);
let output = output.pty_data.unwrap(); assert_eq!( output.data.as_ref(), b"error: debug shell failed: debug shell request is missing failed step context\r\n" );
let exit = exit.exec_exit.unwrap(); assert_eq!(exit.exit_code, 1); assert_eq!(exit.error, detail); }
#[test] fn debug_shell_uses_failed_step_environment_and_devshell() { let open = v1::OpenDebugShell { term: "xterm".to_owned(), failed_step: Some(v1::ExecStart { argv: vec!["/run/current-system/sw/bin/bash".to_owned()], env: vec![ "HOME=/workspace".to_owned(), "CI=true".to_owned(), "SECRET=value".to_owned(), ], cwd: "/workspace/repo".to_owned(), user: "65534".to_owned(), timeout_seconds: 60, stdio_vsock_port: 0, }), shell_prelude: ". /run/spindle/devshell-env.sh; ".to_owned(), ..Default::default() };
let user = exec::ResolvedUser { name: "spindle-workflow".to_owned(), uid: 1000, gid: 1000, home: OsString::from("/workspace"), shell: OsString::from("/bin/bash"), };
let failed_step = open.failed_step.as_ref().unwrap(); let spec = debug_shell_spec_for_user(&open, failed_step, &user); let env: HashMap<_, _> = spec.env.into_iter().collect();
assert_eq!(spec.program, "/bin/bash"); assert_eq!(spec.cwd, Some(Path::new("/workspace/repo").to_path_buf())); assert_eq!( env.get(OsStr::new("XDG_CACHE_HOME")).unwrap(), "/workspace/.cache" ); assert_eq!(env.get(OsStr::new("CI")).unwrap(), "true"); assert_eq!(env.get(OsStr::new("SECRET")).unwrap(), "value"); assert_eq!(env.get(OsStr::new("TERM")).unwrap(), "xterm"); assert_eq!(spec.timeout, None); assert_eq!(spec.args[0], "-lc"); assert_eq!(spec.args[2], "/bin/bash"); assert!( spec.args[1] .to_string_lossy() .contains(". /run/spindle/devshell-env.sh") ); }}