This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / shuttle / src / exec.rs
8.6 kB 284 lines
1use crate::command::{self, OutKind, Spec}; 2use crate::protocol::{self, Message, v1}; 3use nix::unistd::{Group, User}; 4use std::ffi::OsString; 5use std::path::PathBuf; 6use std::time::Duration; 7use tokio::io::AsyncWriteExt; 8use tokio::sync::mpsc::Sender; 9use tokio_vsock::{VsockAddr, VsockStream}; 10use tracing::{info, warn}; 11 12const DEFAULT_USER: &str = "spindle-workflow"; 13const VSOCK_CONNECT_TIMEOUT: Duration = Duration::from_secs(10); 14 15pub async fn run(id: String, req: v1::ExecStart, out: Sender<Message>, host_cid: u32) { 16 let send_exit = async |exit_code: i32, error: Option<String>, timed_out: bool| { 17 let msg = Message { 18 id: id.clone(), 19 exec_exit: Some(v1::ExecExit { 20 exit_code, 21 error: protocol::error_or_empty(error), 22 timed_out, 23 }), 24 ..Default::default() 25 }; 26 let _ = out.send(msg).await; 27 }; 28 29 if req.argv.is_empty() { 30 send_exit(127, Some("missing argv".to_owned()), false).await; 31 return; 32 } 33 34 let user = if req.user.is_empty() { 35 DEFAULT_USER 36 } else { 37 req.user.as_str() 38 }; 39 let run_as = match resolve_user(user) { 40 Ok(run_as) => run_as, 41 Err(err) => { 42 send_exit(127, Some(err), false).await; 43 return; 44 } 45 }; 46 47 let mut env = run_as.login_env(); 48 let runtime_dir = run_as.runtime_dir(); 49 match runtime_dir.try_exists() { 50 Ok(true) => env.push(( 51 OsString::from("XDG_RUNTIME_DIR"), 52 runtime_dir.into_os_string(), 53 )), 54 Ok(false) => {} 55 Err(err) => warn!(error = %err, "could not stat XDG_RUNTIME_DIR for workflow user"), 56 } 57 58 let mut spec = Spec::new(req.argv[0].clone()) 59 .args(req.argv[1..].iter().cloned()) 60 .envs(env) 61 .envs(parse_env(&req.env)) 62 .run_as(run_as.uid, run_as.gid); 63 64 if req.stdio_vsock_port == 0 { 65 send_exit( 66 1, 67 Some("exec is missing a stdio vsock port".to_owned()), 68 false, 69 ) 70 .await; 71 return; 72 } 73 let addr = VsockAddr::new(host_cid, req.stdio_vsock_port); 74 let conn = match tokio::time::timeout(VSOCK_CONNECT_TIMEOUT, VsockStream::connect(addr)).await { 75 Ok(Ok(conn)) => conn, 76 Ok(Err(error)) => { 77 send_exit(1, Some(format!("dial host stdio port: {error}")), false).await; 78 return; 79 } 80 Err(_) => { 81 send_exit(1, Some("dial host stdio port: timed out".to_owned()), false).await; 82 return; 83 } 84 }; 85 spec = spec.piped_stdin(); 86 if !req.cwd.is_empty() { 87 spec = spec.cwd(req.cwd.clone()); 88 } 89 let timeout = 90 (req.timeout_seconds > 0).then(|| Duration::from_secs(u64::from(req.timeout_seconds))); 91 if let Some(timeout) = timeout { 92 spec = spec.timeout(timeout); 93 } 94 95 info!( 96 %id, 97 user = %run_as.name, 98 uid = run_as.uid, 99 gid = run_as.gid, 100 argv = ?req.argv, 101 cwd = ?req.cwd, 102 "starting exec" 103 ); 104 105 let mut cmd = match command::spawn_streaming(spec) { 106 Ok(cmd) => cmd, 107 Err(err) => { 108 send_exit(127, Some(err.to_string()), false).await; 109 return; 110 } 111 }; 112 113 let (mut conn_reader, mut conn_writer) = tokio::io::split(conn); 114 // dropping this also closes the socket if the exec gets cancelled 115 let mut child_stdin = cmd.stdin.take().expect("piped_stdin"); 116 let _stdin_pump = AbortOnDrop(tokio::spawn(async move { 117 let _ = tokio::io::copy(&mut conn_reader, &mut child_stdin).await; 118 })); 119 120 let (mut events, exit_task) = cmd.into_parts(); 121 let mut stdio_error = None; 122 while let Some(event) = events.recv().await { 123 match event.kind { 124 OutKind::Stdout => { 125 if let Err(error) = conn_writer.write_all(&event.data).await { 126 stdio_error = Some(format!("forward stdout to host: {error}")); 127 // stop these before their pipes fill and block the child 128 events.close(); 129 break; 130 } 131 } 132 OutKind::Stderr => { 133 let _ = out 134 .send(Message { 135 id: id.clone(), 136 exec_stderr: Some(v1::ExecStderr { 137 data: event.data.into(), 138 }), 139 ..Default::default() 140 }) 141 .await; 142 } 143 } 144 } 145 // the child might keep reading stdin after it closes stdout 146 let exit = match exit_task 147 .await 148 .unwrap_or_else(|error| Err(anyhow::anyhow!("command supervisor failed: {error}"))) 149 { 150 Ok(exit) => exit, 151 Err(err) => { 152 send_exit(127, Some(err.to_string()), false).await; 153 return; 154 } 155 }; 156 157 // losing stdout fails the exec even if the child exited cleanly 158 if let Some(error) = stdio_error { 159 send_exit(1, Some(error), false).await; 160 return; 161 } 162 send_exit(exit.exit_code, exit.error, exit.timed_out).await 163} 164 165// aborts the task when its exec goes away 166struct AbortOnDrop(tokio::task::JoinHandle<()>); 167 168impl Drop for AbortOnDrop { 169 fn drop(&mut self) { 170 self.0.abort(); 171 } 172} 173 174#[derive(Clone, Debug)] 175struct ResolvedUser { 176 name: String, 177 uid: u32, 178 gid: u32, 179 home: OsString, 180 shell: OsString, 181} 182 183impl ResolvedUser { 184 fn login_env(&self) -> Vec<(OsString, OsString)> { 185 let xdg_cache_home = PathBuf::from(&self.home).join(".cache"); 186 vec![ 187 (OsString::from("USER"), OsString::from(&self.name)), 188 (OsString::from("LOGNAME"), OsString::from(&self.name)), 189 (OsString::from("HOME"), self.home.clone()), 190 ( 191 OsString::from("XDG_CACHE_HOME"), 192 xdg_cache_home.into_os_string(), 193 ), 194 (OsString::from("SHELL"), self.shell.clone()), 195 ] 196 } 197 198 fn runtime_dir(&self) -> PathBuf { 199 PathBuf::from(format!("/run/user/{}", self.uid)) 200 } 201} 202 203fn resolve_user(spec: &str) -> Result<ResolvedUser, String> { 204 let spec = spec.trim(); 205 if spec.is_empty() { 206 return resolve_user(DEFAULT_USER); 207 } 208 209 let (user_part, group_part) = spec 210 .split_once(':') 211 .map(|(user, group)| (user, Some(group))) 212 .unwrap_or((spec, None)); 213 214 let mut user = lookup_user(user_part)?; 215 if let Some(group) = group_part.filter(|group| !group.is_empty()) { 216 user.gid = lookup_group(group)?; 217 } 218 219 if user.uid == 0 || user.gid == 0 { 220 return Err(format!("refusing to run exec as privileged user {spec:?}")); 221 } 222 223 Ok(user) 224} 225 226fn lookup_user(name: &str) -> Result<ResolvedUser, String> { 227 match User::from_name(name) { 228 Ok(Some(user)) => Ok(ResolvedUser { 229 name: name.to_owned(), 230 uid: user.uid.as_raw(), 231 gid: user.gid.as_raw(), 232 home: user.dir.into_os_string(), 233 shell: user.shell.into_os_string(), 234 }), 235 Ok(None) => { 236 let uid = name 237 .parse::<u32>() 238 .map_err(|_| format!("workflow user {name:?} was not found"))?; 239 Ok(ResolvedUser { 240 name: name.to_owned(), 241 uid, 242 gid: uid, 243 home: OsString::from("/"), 244 shell: OsString::from("/bin/sh"), 245 }) 246 } 247 Err(error) => Err(format!("lookup workflow user {name:?}: {error}")), 248 } 249} 250 251fn lookup_group(name: &str) -> Result<u32, String> { 252 match Group::from_name(name) { 253 Ok(Some(group)) => Ok(group.gid.as_raw()), 254 Ok(None) => name 255 .parse::<u32>() 256 .map_err(|_| format!("workflow group {name:?} was not found")), 257 Err(error) => Err(format!("lookup workflow group {name:?}: {error}")), 258 } 259} 260 261fn parse_env(values: &[String]) -> Vec<(OsString, OsString)> { 262 values 263 .iter() 264 .filter_map(|value| value.split_once('=')) 265 .map(|(key, value)| (OsString::from(key), OsString::from(value))) 266 .collect() 267} 268 269#[cfg(test)] 270mod tests { 271 use super::*; 272 273 #[test] 274 fn refuses_root_exec_user() { 275 let err = resolve_user("root").unwrap_err(); 276 assert!(err.contains("refusing to run exec as privileged user")); 277 } 278 279 #[test] 280 fn refuses_root_exec_group() { 281 let err = resolve_user("65534:0").unwrap_err(); 282 assert!(err.contains("refusing to run exec as privileged user")); 283 } 284}