This repository has no description
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}