Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296use crate::command::{self, OutKind, Spec};use crate::protocol::{self, Message, v1};use nix::unistd::{Group, User};use std::ffi::OsString;use std::path::PathBuf;use std::time::Duration;use tokio::io::AsyncWriteExt;use tokio::sync::mpsc::Sender;use tokio_vsock::{VsockAddr, VsockStream};use tracing::{info, warn};
const DEFAULT_USER: &str = "spindle-workflow";const VSOCK_CONNECT_TIMEOUT: Duration = Duration::from_secs(10);
pub async fn run(id: String, req: v1::ExecStart, out: Sender<Message>, host_cid: u32) { let send_exit = async |exit_code: i32, error: Option<String>, timed_out: bool| { let msg = Message { id: id.clone(), exec_exit: Some(v1::ExecExit { exit_code, error: protocol::error_or_empty(error), timed_out, }), ..Default::default() }; let _ = out.send(msg).await; };
let (mut spec, run_as) = match request_spec(&req) { Ok(result) => result, Err(err) => { send_exit(127, Some(err), false).await; return; } };
if req.stdio_vsock_port == 0 { send_exit( 1, Some("exec is missing a stdio vsock port".to_owned()), false, ) .await; return; } let addr = VsockAddr::new(host_cid, req.stdio_vsock_port); let conn = match tokio::time::timeout(VSOCK_CONNECT_TIMEOUT, VsockStream::connect(addr)).await { Ok(Ok(conn)) => conn, Ok(Err(error)) => { send_exit(1, Some(format!("dial host stdio port: {error}")), false).await; return; } Err(_) => { send_exit(1, Some("dial host stdio port: timed out".to_owned()), false).await; return; } }; spec = spec.piped_stdin();
info!( %id, user = %run_as.name, uid = run_as.uid, gid = run_as.gid, argv = ?req.argv, cwd = ?req.cwd, "starting exec" );
let mut cmd = match command::spawn_streaming(spec) { Ok(cmd) => cmd, Err(err) => { send_exit(127, Some(err.to_string()), false).await; return; } };
let (mut conn_reader, mut conn_writer) = tokio::io::split(conn); // dropping this also closes the socket if the exec gets cancelled let mut child_stdin = cmd.stdin.take().expect("piped_stdin"); let _stdin_pump = AbortOnDrop(tokio::spawn(async move { let _ = tokio::io::copy(&mut conn_reader, &mut child_stdin).await; }));
let (mut events, exit_task) = cmd.into_parts(); let mut stdio_error = None; while let Some(event) = events.recv().await { match event.kind { OutKind::Stdout => { if let Err(error) = conn_writer.write_all(&event.data).await { stdio_error = Some(format!("forward stdout to host: {error}")); // stop these before their pipes fill and block the child events.close(); break; } } OutKind::Stderr => { let _ = out .send(Message { id: id.clone(), exec_stderr: Some(v1::ExecStderr { data: event.data.into(), }), ..Default::default() }) .await; } } } // the child might keep reading stdin after it closes stdout let exit = match exit_task .await .unwrap_or_else(|error| Err(anyhow::anyhow!("command supervisor failed: {error}"))) { Ok(exit) => exit, Err(err) => { send_exit(127, Some(err.to_string()), false).await; return; } };
// losing stdout fails the exec even if the child exited cleanly if let Some(error) = stdio_error { send_exit(1, Some(error), false).await; return; } send_exit(exit.exit_code, exit.error, exit.timed_out).await}
// aborts the task when its exec goes awaystruct AbortOnDrop(tokio::task::JoinHandle<()>);
impl Drop for AbortOnDrop { fn drop(&mut self) { self.0.abort(); }}
pub(crate) fn request_spec(req: &v1::ExecStart) -> Result<(Spec, ResolvedUser), String> { if req.argv.is_empty() { return Err("missing argv".to_owned()); } let run_as = resolve_user(&req.user)?; let spec = request_spec_for_user(req, &run_as); Ok((spec, run_as))}
pub(crate) fn request_spec_for_user(req: &v1::ExecStart, run_as: &ResolvedUser) -> Spec { let mut env = run_as.login_env(); let runtime_dir = run_as.runtime_dir(); match runtime_dir.try_exists() { Ok(true) => env.push(( OsString::from("XDG_RUNTIME_DIR"), runtime_dir.into_os_string(), )), Ok(false) => {} Err(err) => warn!(error = %err, "could not stat XDG_RUNTIME_DIR for workflow user"), }
let mut spec = Spec::new(req.argv[0].clone()) .args(req.argv[1..].iter().cloned()) .envs(env) .envs(parse_env(&req.env)) .run_as(run_as.uid, run_as.gid); if !req.cwd.is_empty() { spec = spec.cwd(req.cwd.clone()); } let timeout = (req.timeout_seconds > 0).then(|| Duration::from_secs(u64::from(req.timeout_seconds))); if let Some(timeout) = timeout { spec = spec.timeout(timeout); } spec}
#[derive(Clone, Debug)]pub(crate) struct ResolvedUser { pub(crate) name: String, pub(crate) uid: u32, pub(crate) gid: u32, pub(crate) home: OsString, pub(crate) shell: OsString,}
impl ResolvedUser { fn login_env(&self) -> Vec<(OsString, OsString)> { let xdg_cache_home = PathBuf::from(&self.home).join(".cache"); vec![ (OsString::from("USER"), OsString::from(&self.name)), (OsString::from("LOGNAME"), OsString::from(&self.name)), (OsString::from("HOME"), self.home.clone()), ( OsString::from("XDG_CACHE_HOME"), xdg_cache_home.into_os_string(), ), (OsString::from("SHELL"), self.shell.clone()), ] }
fn runtime_dir(&self) -> PathBuf { PathBuf::from(format!("/run/user/{}", self.uid)) }
pub(crate) fn shell(&self) -> &OsString { &self.shell }}
pub(crate) fn resolve_user(spec: &str) -> Result<ResolvedUser, String> { let spec = spec.trim(); if spec.is_empty() { return resolve_user(DEFAULT_USER); }
let (user_part, group_part) = spec .split_once(':') .map(|(user, group)| (user, Some(group))) .unwrap_or((spec, None));
let mut user = lookup_user(user_part)?; if let Some(group) = group_part.filter(|group| !group.is_empty()) { user.gid = lookup_group(group)?; }
if user.uid == 0 || user.gid == 0 { return Err(format!("refusing to run exec as privileged user {spec:?}")); }
Ok(user)}
fn lookup_user(name: &str) -> Result<ResolvedUser, String> { match name.parse::<u32>() { Ok(uid) => Ok(ResolvedUser { name: name.to_owned(), uid, gid: uid, home: OsString::from("/"), shell: OsString::from("/bin/sh"), }), Err(_) => match User::from_name(name) { Ok(Some(user)) => Ok(ResolvedUser { name: name.to_owned(), uid: user.uid.as_raw(), gid: user.gid.as_raw(), home: user.dir.into_os_string(), shell: user.shell.into_os_string(), }), Ok(None) => Err(format!("workflow user {name:?} was not found")), Err(error) => Err(format!("lookup workflow user {name:?}: {error}")), }, }}
fn lookup_group(name: &str) -> Result<u32, String> { match name.parse::<u32>() { Ok(gid) => Ok(gid), Err(_) => match Group::from_name(name) { Ok(Some(group)) => Ok(group.gid.as_raw()), Ok(None) => Err(format!("workflow group {name:?} was not found")), Err(error) => Err(format!("lookup workflow group {name:?}: {error}")), }, }}
fn parse_env(values: &[String]) -> Vec<(OsString, OsString)> { values .iter() .filter_map(|value| value.split_once('=')) .map(|(key, value)| (OsString::from(key), OsString::from(value))) .collect()}
#[cfg(test)]mod tests { use super::*;
#[test] fn refuses_exec_as_uid_zero() { let err = resolve_user("0").unwrap_err(); assert!(err.contains("refusing to run exec as privileged user")); }
#[test] fn refuses_exec_as_gid_zero() { let err = resolve_user("65534:0").unwrap_err(); assert!(err.contains("refusing to run exec as privileged user")); }
#[test] fn resolves_numeric_spec_without_a_user_database() { let user = resolve_user("65534:65533").unwrap(); assert_eq!((user.uid, user.gid), (65534, 65533)); }}