Skip to main content

isb_core/
exec.rs

1//! Running commands in a sandbox over the incus exec websocket API.
2//!
3//! - argv is passed as a list, never joined into `sh -c`.
4//! - Output is streamed as produced, and there is no default timeout.
5//! - stdin, when the caller asks for it, is forwarded and closed properly. A
6//!   caller whose own stdin is an inherited pipe that never reaches EOF does not
7//!   hang: the forwarder polls and stops when the command exits.
8//! - Closing the control websocket kills the command in incus, so it is held open
9//!   until the operation finishes.
10
11use std::collections::BTreeMap;
12use std::io::Write;
13use std::os::unix::net::UnixStream;
14use std::sync::atomic::{AtomicBool, Ordering};
15use std::sync::mpsc::{Receiver, SyncSender, sync_channel};
16use std::sync::{Arc, Mutex};
17use std::thread::JoinHandle;
18use std::time::{Duration, Instant};
19
20use serde_json::{Value, json};
21use tungstenite::{Message, WebSocket};
22
23use crate::client::{Client, Reply, encode_segment};
24use crate::error::{Error, Result};
25
26type Ws = WebSocket<UnixStream>;
27
28/// Output chunks a command may run ahead of its reader.
29const EVENTS_QUEUED: usize = 256;
30/// Input chunks queued for a command's stdin.
31const STDIN_QUEUED: usize = 16;
32
33/// Where the command's stdin comes from.
34#[derive(Debug, Clone, Default)]
35pub enum Stdin {
36    /// Closed immediately (like `/dev/null`).
37    #[default]
38    Null,
39    /// These bytes, then EOF.
40    Bytes(Vec<u8>),
41    /// This process's own stdin (fd 0), until EOF or until the command exits.
42    Inherit,
43    /// Written through [`ExecStream::write_stdin`], ended by
44    /// [`ExecStream::close_stdin`].
45    Piped,
46}
47
48/// Per-call exec options. Unset fields fall back to the sandbox's exec defaults.
49#[derive(Debug, Clone, Default)]
50pub struct ExecOptions {
51    pub cwd: Option<String>,
52    /// A guest user name, `uid`, or `uid:gid`.
53    pub user: Option<String>,
54    pub env: BTreeMap<String, String>,
55    /// Run through the user's login shell.
56    pub login: Option<bool>,
57    /// Allocate a pseudo-terminal (stdout and stderr are merged, as on any tty).
58    pub tty: bool,
59    pub width: Option<u16>,
60    pub height: Option<u16>,
61    /// Kill the command after this long. None (the default) means no limit.
62    pub timeout: Option<Duration>,
63    pub stdin: Stdin,
64}
65
66impl ExecOptions {
67    pub fn cwd(mut self, c: impl Into<String>) -> Self {
68        self.cwd = Some(c.into());
69        self
70    }
71    pub fn user(mut self, u: impl Into<String>) -> Self {
72        self.user = Some(u.into());
73        self
74    }
75    pub fn env(mut self, k: impl Into<String>, v: impl Into<String>) -> Self {
76        self.env.insert(k.into(), v.into());
77        self
78    }
79    pub fn login(mut self, l: bool) -> Self {
80        self.login = Some(l);
81        self
82    }
83    pub fn tty(mut self, t: bool) -> Self {
84        self.tty = t;
85        self
86    }
87    pub fn timeout(mut self, t: Duration) -> Self {
88        self.timeout = Some(t);
89        self
90    }
91    pub fn stdin(mut self, s: Stdin) -> Self {
92        self.stdin = s;
93        self
94    }
95}
96
97/// Captured result of [`crate::Sandbox::exec`].
98#[derive(Debug, Clone, Default, PartialEq, Eq)]
99pub struct ExecOutput {
100    pub exit_code: i32,
101    pub stdout: Vec<u8>,
102    pub stderr: Vec<u8>,
103}
104
105impl ExecOutput {
106    pub fn success(&self) -> bool {
107        self.exit_code == 0
108    }
109    pub fn stdout_text(&self) -> String {
110        String::from_utf8_lossy(&self.stdout).into_owned()
111    }
112    pub fn stderr_text(&self) -> String {
113        String::from_utf8_lossy(&self.stderr).into_owned()
114    }
115}
116
117/// A chunk of output, in the order it was produced per stream.
118#[derive(Debug, Clone, PartialEq, Eq)]
119pub enum ExecEvent {
120    Stdout(Vec<u8>),
121    Stderr(Vec<u8>),
122}
123
124/// A guest user resolved to ids.
125#[derive(Debug, Clone, Default, PartialEq, Eq)]
126pub struct GuestUser {
127    pub name: Option<String>,
128    pub uid: u32,
129    pub gid: u32,
130    pub home: Option<String>,
131    pub shell: Option<String>,
132}
133
134/// Parse a `getent passwd` line.
135pub fn parse_passwd(line: &str) -> Option<GuestUser> {
136    let f: Vec<&str> = line.trim().split(':').collect();
137    if f.len() < 7 {
138        return None;
139    }
140    Some(GuestUser {
141        name: Some(f[0].to_string()),
142        uid: f[2].parse().ok()?,
143        gid: f[3].parse().ok()?,
144        home: Some(f[5].to_string()).filter(|s| !s.is_empty()),
145        shell: Some(f[6].to_string()).filter(|s| !s.is_empty()),
146    })
147}
148
149mod user;
150pub use user::resolve_user;
151
152/// A fully resolved exec request.
153#[derive(Debug, Clone, Default)]
154pub(crate) struct Request {
155    pub argv: Vec<String>,
156    pub cwd: Option<String>,
157    pub uid: u32,
158    pub gid: u32,
159    pub env: BTreeMap<String, String>,
160    pub tty: bool,
161    pub width: Option<u16>,
162    pub height: Option<u16>,
163}
164
165impl Request {
166    fn root() -> Self {
167        Request::default()
168    }
169}
170
171/// Merge sandbox defaults with call options and resolve the user.
172pub(crate) fn build_request(
173    client: &Client,
174    instance: &str,
175    argv: &[String],
176    defaults: &crate::spec::ExecDefaults,
177    opts: &ExecOptions,
178) -> Result<Request> {
179    if argv.is_empty() {
180        return Err(Error::invalid("exec needs a command"));
181    }
182    let mut env = defaults.env.clone();
183    env.extend(opts.env.clone());
184    let user = opts.user.clone().or_else(|| defaults.user.clone());
185    let login = opts.login.unwrap_or(defaults.login);
186    let gu = match &user {
187        Some(u) => Some(resolve_user(client, instance, u)?),
188        None => None,
189    };
190    if let Some(g) = &gu {
191        if let Some(h) = &g.home {
192            env.entry("HOME".into()).or_insert_with(|| h.clone());
193        }
194        if let Some(n) = &g.name {
195            env.entry("USER".into()).or_insert_with(|| n.clone());
196            env.entry("LOGNAME".into()).or_insert_with(|| n.clone());
197        }
198    }
199    let mut final_argv = argv.to_vec();
200    if login {
201        let shell = gu
202            .as_ref()
203            .and_then(|g| g.shell.clone())
204            .filter(|s| !s.ends_with("nologin") && !s.ends_with("/false"))
205            .unwrap_or_else(|| "/bin/sh".into());
206        // argv stays a list: the shell only runs `exec "$@"`.
207        let mut v = vec![
208            shell,
209            "-l".into(),
210            "-c".into(),
211            "exec \"$@\"".into(),
212            "isb".into(),
213        ];
214        v.extend(final_argv);
215        final_argv = v;
216    }
217    if opts.tty {
218        env.entry("TERM".into())
219            .or_insert_with(|| std::env::var("TERM").unwrap_or_else(|_| "xterm-256color".into()));
220    }
221    Ok(Request {
222        argv: final_argv,
223        cwd: opts
224            .cwd
225            .clone()
226            .or_else(|| defaults.cwd.clone())
227            .or_else(|| gu.as_ref().and_then(|g| g.home.clone())),
228        uid: gu.as_ref().map(|g| g.uid).unwrap_or(0),
229        gid: gu.as_ref().map(|g| g.gid).unwrap_or(0),
230        env,
231        tty: opts.tty,
232        width: opts.width,
233        height: opts.height,
234    })
235}
236
237/// A running command. Iterate it for output; [`ExecStream::wait`] for the exit code.
238pub struct ExecStream {
239    events: Receiver<ExecEvent>,
240    control: Arc<Mutex<Option<Ws>>>,
241    stdin_tx: Option<SyncSender<Option<Vec<u8>>>>,
242    waiter: Option<JoinHandle<Result<i32>>>,
243    stop: Arc<AtomicBool>,
244    tty: bool,
245}
246
247/// A cloneable handle for driving a running command (stdin, signals, window
248/// size) from other threads while one thread reads its output.
249#[derive(Clone)]
250pub struct ExecController {
251    control: Arc<Mutex<Option<Ws>>>,
252    stdin_tx: Option<SyncSender<Option<Vec<u8>>>>,
253    tty: bool,
254}
255
256impl ExecController {
257    /// Write to stdin (only with [`Stdin::Piped`]).
258    pub fn write_stdin(&self, data: &[u8]) -> Result<()> {
259        match &self.stdin_tx {
260            Some(tx) => tx
261                .send(Some(data.to_vec()))
262                .map_err(|_| Error::WebSocket("stdin is closed".into())),
263            None => Err(Error::invalid("stdin is not piped")),
264        }
265    }
266
267    /// Send EOF on stdin (only with [`Stdin::Piped`]). Later writes fail.
268    pub fn close_stdin(&self) -> Result<()> {
269        match &self.stdin_tx {
270            Some(tx) => {
271                let _ = tx.send(None);
272                Ok(())
273            }
274            None => Err(Error::invalid("stdin is not piped")),
275        }
276    }
277
278    pub fn signal(&self, signal: i32) -> Result<()> {
279        send_control(
280            &self.control,
281            json!({"command": "signal", "signal": signal}),
282        )
283    }
284
285    pub fn resize(&self, width: u16, height: u16) -> Result<()> {
286        if !self.tty {
287            return Ok(());
288        }
289        send_control(
290            &self.control,
291            json!({"command": "window-resize", "args": {"width": width.to_string(), "height": height.to_string()}}),
292        )
293    }
294}
295
296impl ExecStream {
297    /// A handle for stdin, signals and resizes usable from other threads.
298    pub fn controller(&self) -> ExecController {
299        ExecController {
300            control: self.control.clone(),
301            stdin_tx: self.stdin_tx.clone(),
302            tty: self.tty,
303        }
304    }
305
306    /// Next chunk of output, blocking. `None` once all output has been read.
307    pub fn next_event(&mut self) -> Option<ExecEvent> {
308        self.events.recv().ok()
309    }
310
311    /// Next chunk of output within `wait`: `Ok(None)` if none came, `Err`
312    /// with the exit code once the command ended and its output was read.
313    pub fn poll_event(
314        &mut self,
315        wait: Duration,
316    ) -> std::result::Result<Option<ExecEvent>, Result<i32>> {
317        use std::sync::mpsc::RecvTimeoutError;
318        let r = if wait.is_zero() {
319            self.events.try_recv().map_err(|e| match e {
320                std::sync::mpsc::TryRecvError::Empty => RecvTimeoutError::Timeout,
321                std::sync::mpsc::TryRecvError::Disconnected => RecvTimeoutError::Disconnected,
322            })
323        } else {
324            self.events.recv_timeout(wait)
325        };
326        match r {
327            Ok(ev) => Ok(Some(ev)),
328            Err(RecvTimeoutError::Timeout) => Ok(None),
329            Err(RecvTimeoutError::Disconnected) => Err(self.finish()),
330        }
331    }
332
333    /// Write to the command's stdin (only with [`Stdin::Piped`]).
334    pub fn write_stdin(&self, data: &[u8]) -> Result<()> {
335        match &self.stdin_tx {
336            Some(tx) => tx
337                .send(Some(data.to_vec()))
338                .map_err(|_| Error::WebSocket("stdin is closed".into())),
339            None => Err(Error::invalid("stdin is not piped")),
340        }
341    }
342
343    /// Send EOF on stdin (only with [`Stdin::Piped`]).
344    pub fn close_stdin(&mut self) {
345        if let Some(tx) = self.stdin_tx.take() {
346            let _ = tx.send(None);
347        }
348    }
349
350    /// Send a signal to the command (e.g. 2 for SIGINT, 15 for SIGTERM).
351    pub fn signal(&self, signal: i32) -> Result<()> {
352        send_control(
353            &self.control,
354            json!({"command": "signal", "signal": signal}),
355        )
356    }
357
358    /// Resize the pseudo-terminal (tty mode only).
359    pub fn resize(&self, width: u16, height: u16) -> Result<()> {
360        if !self.tty {
361            return Ok(());
362        }
363        send_control(
364            &self.control,
365            json!({"command": "window-resize", "args": {"width": width.to_string(), "height": height.to_string()}}),
366        )
367    }
368
369    /// Discard any remaining output and return the exit code.
370    pub fn wait(mut self) -> Result<i32> {
371        while self.events.recv().is_ok() {}
372        self.finish()
373    }
374
375    fn finish(&mut self) -> Result<i32> {
376        let r = match self.waiter.take() {
377            Some(h) => h
378                .join()
379                .unwrap_or_else(|_| Err(Error::Protocol("exec waiter panicked".into()))),
380            None => Err(Error::Protocol("exec already finished".into())),
381        };
382        self.stop.store(true, Ordering::SeqCst);
383        r
384    }
385
386    /// Collect everything into an [`ExecOutput`].
387    pub fn collect_output(mut self) -> Result<ExecOutput> {
388        let mut out = ExecOutput::default();
389        while let Some(ev) = self.next_event() {
390            match ev {
391                ExecEvent::Stdout(b) => out.stdout.extend(b),
392                ExecEvent::Stderr(b) => out.stderr.extend(b),
393            }
394        }
395        out.exit_code = self.finish()?;
396        Ok(out)
397    }
398}
399
400impl Iterator for ExecStream {
401    type Item = ExecEvent;
402    fn next(&mut self) -> Option<ExecEvent> {
403        self.next_event()
404    }
405}
406
407impl Drop for ExecStream {
408    fn drop(&mut self) {
409        self.stop.store(true, Ordering::SeqCst);
410    }
411}
412
413fn send_control(control: &Arc<Mutex<Option<Ws>>>, msg: Value) -> Result<()> {
414    let mut g = control.lock().unwrap_or_else(|p| p.into_inner());
415    match g.as_mut() {
416        Some(ws) => ws
417            .send(Message::text(msg.to_string()))
418            .map_err(|e| Error::WebSocket(format!("control: {e}"))),
419        None => Err(Error::WebSocket("command has finished".into())),
420    }
421}
422
423fn is_would_block(e: &tungstenite::Error) -> bool {
424    matches!(e, tungstenite::Error::Io(io) if matches!(io.kind(), std::io::ErrorKind::WouldBlock | std::io::ErrorKind::TimedOut))
425}
426
427/// Read one output websocket to the end, forwarding binary frames.
428fn pump_output(mut ws: Ws, tx: SyncSender<ExecEvent>, stderr: bool) {
429    loop {
430        match ws.read() {
431            Ok(Message::Binary(b)) => {
432                if b.is_empty() {
433                    continue;
434                }
435                let ev = if stderr {
436                    ExecEvent::Stderr(b.to_vec())
437                } else {
438                    ExecEvent::Stdout(b.to_vec())
439                };
440                if tx.send(ev).is_err() {
441                    // Receiver gone; keep draining so the server is never blocked.
442                }
443            }
444            // incus ends a stream with an empty text frame (its "barrier"), sent
445            // after the last byte of output.
446            Ok(Message::Text(_)) | Ok(Message::Close(_)) => {
447                let _ = ws.close(None);
448                let _ = ws.flush();
449                break;
450            }
451            Ok(_) => {}
452            Err(_) => break,
453        }
454    }
455}
456
457#[expect(
458    clippy::too_many_lines,
459    reason = "predates the lint ratchet; split it when next changed"
460)]
461pub(crate) fn start_with_timeout(
462    client: &Client,
463    instance: &str,
464    req: Request,
465    stdin: Stdin,
466    timeout: Option<Duration>,
467) -> Result<ExecStream> {
468    let mut body = json!({
469        "command": req.argv,
470        "environment": req.env,
471        "wait-for-websocket": true,
472        "interactive": req.tty,
473        "user": req.uid,
474        "group": req.gid,
475    });
476    if let Some(c) = &req.cwd {
477        body["cwd"] = json!(c);
478    }
479    if req.tty {
480        body["width"] = json!(req.width.unwrap_or(80));
481        body["height"] = json!(req.height.unwrap_or(24));
482    }
483    let path = format!("/1.0/instances/{}/exec", encode_segment(instance));
484    let (operation, meta) =
485        match client.request("POST", &path, Some(&body), client.timeouts.request)? {
486            Reply::Async {
487                operation,
488                metadata,
489            } => (operation, metadata),
490            Reply::Sync(_) => {
491                return Err(Error::Protocol("exec did not return an operation".into()));
492            }
493        };
494    let fds = meta
495        .pointer("/metadata/fds")
496        .and_then(Value::as_object)
497        .cloned()
498        .ok_or_else(|| Error::Protocol("exec operation has no websocket secrets".into()))?;
499    let secret = |k: &str| -> Result<String> {
500        fds.get(k)
501            .and_then(Value::as_str)
502            .map(String::from)
503            .ok_or_else(|| Error::Protocol(format!("exec operation has no {k} websocket")))
504    };
505
506    let control_ws = client.websocket(&operation, &secret("control")?)?;
507    let control = Arc::new(Mutex::new(Some(control_ws)));
508    // Bounded both ways, so a slow reader (a backup streaming to S3) or a
509    // fast writer (a restore) holds the command back instead of buffering
510    // its whole output or input in memory.
511    let (tx, rx) = sync_channel::<ExecEvent>(EVENTS_QUEUED);
512    let stop = Arc::new(AtomicBool::new(false));
513    let mut readers: Vec<JoinHandle<()>> = Vec::new();
514    let mut stdin_tx = None;
515
516    // The stdin source, as a channel of chunks (None = EOF), fed by a thread.
517    let (in_tx, in_rx) = sync_channel::<Option<Vec<u8>>>(STDIN_QUEUED);
518    match stdin {
519        Stdin::Null => {
520            let _ = in_tx.send(None);
521        }
522        Stdin::Bytes(b) => {
523            std::thread::spawn(move || {
524                for c in b.chunks(64 * 1024) {
525                    if in_tx.send(Some(c.to_vec())).is_err() {
526                        return;
527                    }
528                }
529                let _ = in_tx.send(None);
530            });
531        }
532        Stdin::Piped => stdin_tx = Some(in_tx),
533        Stdin::Inherit => {
534            let stop2 = stop.clone();
535            std::thread::spawn(move || forward_fd0(in_tx, stop2));
536        }
537    }
538
539    if req.tty {
540        // One bidirectional socket. Reads poll with a short timeout so the writer
541        // can take the lock between them.
542        let ws = client.websocket(&operation, &secret("0")?)?;
543        ws.get_ref()
544            .set_read_timeout(Some(Duration::from_millis(20)))?;
545        let ws = Arc::new(Mutex::new(ws));
546        let reader_ws = ws.clone();
547        let tx2 = tx.clone();
548        readers.push(std::thread::spawn(move || {
549            loop {
550                let r = {
551                    let mut g = reader_ws.lock().unwrap_or_else(|p| p.into_inner());
552                    g.read()
553                };
554                match r {
555                    Ok(Message::Binary(b)) => {
556                        let _ = tx2.send(ExecEvent::Stdout(b.to_vec()));
557                    }
558                    // End of output (the barrier). incus finishes the operation
559                    // only once the client closes this socket, so close it now;
560                    // waiting for our own stdin to end would hang on a terminal.
561                    Ok(Message::Text(_)) | Ok(Message::Close(_)) => {
562                        let mut g = reader_ws.lock().unwrap_or_else(|p| p.into_inner());
563                        let _ = g.close(None);
564                        let _ = g.flush();
565                        break;
566                    }
567                    Ok(_) => {}
568                    Err(e) if is_would_block(&e) => {
569                        std::thread::sleep(Duration::from_millis(1));
570                    }
571                    Err(_) => break,
572                }
573            }
574        }));
575        let stop2 = stop.clone();
576        std::thread::spawn(move || {
577            while let Ok(Some(chunk)) = in_rx.recv() {
578                if stop2.load(Ordering::SeqCst) {
579                    break;
580                }
581                let mut g = ws.lock().unwrap_or_else(|p| p.into_inner());
582                if g.send(Message::binary(chunk)).is_err() {
583                    break;
584                }
585            }
586        });
587    } else {
588        // Connect all three before anything else: incus starts the command once
589        // stdin, stdout and stderr are attached.
590        let mut ws_in = client.websocket(&operation, &secret("0")?)?;
591        let ws_out = client.websocket(&operation, &secret("1")?)?;
592        let ws_err = client.websocket(&operation, &secret("2")?)?;
593        let tx1 = tx.clone();
594        readers.push(std::thread::spawn(move || pump_output(ws_out, tx1, false)));
595        let tx2 = tx.clone();
596        readers.push(std::thread::spawn(move || pump_output(ws_err, tx2, true)));
597        let stop2 = stop.clone();
598        std::thread::spawn(move || {
599            loop {
600                match in_rx.recv() {
601                    Ok(Some(chunk)) => {
602                        if stop2.load(Ordering::SeqCst)
603                            || ws_in.send(Message::binary(chunk)).is_err()
604                        {
605                            return;
606                        }
607                    }
608                    // EOF (or the source went away): an empty text frame is incus'
609                    // end-of-stdin barrier; then close cleanly.
610                    Ok(None) | Err(_) => {
611                        let _ = ws_in.send(Message::text(""));
612                        let _ = ws_in.close(None);
613                        let _ = ws_in.flush();
614                        return;
615                    }
616                }
617            }
618        });
619    }
620    drop(tx);
621
622    let client2 = client.clone();
623    let control2 = control.clone();
624    let argv_desc = req.argv.join(" ");
625    let stop3 = stop.clone();
626    let waiter = std::thread::spawn(move || -> Result<i32> {
627        let started = Instant::now();
628        let mut killed_at: Option<Instant> = None;
629        let result = loop {
630            let secs = match (timeout, killed_at) {
631                (_, Some(_)) => 1,
632                (Some(t), None) => t.saturating_sub(started.elapsed()).as_secs().min(30),
633                (None, None) => 30,
634            };
635            if secs == 0 {
636                std::thread::sleep(Duration::from_millis(50));
637            }
638            match client2.poll_operation(&operation, "exec", secs) {
639                Ok(Some(meta)) => {
640                    break Ok(meta.get("return").and_then(Value::as_i64).unwrap_or(-1) as i32);
641                }
642                Ok(None) => {}
643                Err(e) => break Err(e),
644            }
645            if let Some(t) = timeout {
646                match killed_at {
647                    None if started.elapsed() >= t => {
648                        let _ = send_control(&control2, json!({"command": "signal", "signal": 9}));
649                        killed_at = Some(Instant::now());
650                    }
651                    Some(k) if k.elapsed() >= Duration::from_secs(10) => {
652                        break Err(Error::ExecTimeout {
653                            argv: argv_desc.clone(),
654                            timeout: t,
655                        });
656                    }
657                    _ => {}
658                }
659            }
660        };
661        // Output sockets close once the process and its mirrors are done.
662        for r in readers {
663            let _ = r.join();
664        }
665        stop3.store(true, Ordering::SeqCst);
666        // Only now is it safe to drop control: closing it earlier kills the command.
667        if let Some(mut ws) = control2.lock().unwrap_or_else(|p| p.into_inner()).take() {
668            let _ = ws.close(None);
669            let _ = ws.flush();
670        }
671        match (result, killed_at, timeout) {
672            (Ok(_), Some(_), Some(t)) => Err(Error::ExecTimeout {
673                argv: argv_desc,
674                timeout: t,
675            }),
676            (r, _, _) => r,
677        }
678    });
679
680    Ok(ExecStream {
681        events: rx,
682        control,
683        stdin_tx,
684        waiter: Some(waiter),
685        stop,
686        tty: req.tty,
687    })
688}
689
690/// Forward fd 0 until EOF, or until `stop` is set. Polls so that an inherited
691/// stdin that never reaches EOF cannot keep anything alive.
692fn forward_fd0(tx: SyncSender<Option<Vec<u8>>>, stop: Arc<AtomicBool>) {
693    use rustix::event::{PollFd, PollFlags, poll};
694    let stdin = std::io::stdin();
695    let mut buf = vec![0u8; 64 * 1024];
696    loop {
697        if stop.load(Ordering::SeqCst) {
698            return;
699        }
700        let ready = {
701            let mut fds = [PollFd::new(&stdin, PollFlags::IN)];
702            let ts = rustix::event::Timespec {
703                tv_sec: 0,
704                tv_nsec: 100_000_000,
705            };
706            match poll(&mut fds, Some(&ts)) {
707                Ok(n) => n > 0,
708                Err(rustix::io::Errno::INTR) => false,
709                Err(_) => {
710                    let _ = tx.send(None);
711                    return;
712                }
713            }
714        };
715        if !ready {
716            continue;
717        }
718        match rustix::io::read(&stdin, &mut buf) {
719            Ok(0) => {
720                let _ = tx.send(None);
721                return;
722            }
723            Ok(n) => {
724                if tx.send(Some(buf[..n].to_vec())).is_err() {
725                    return;
726                }
727            }
728            Err(rustix::io::Errno::INTR) | Err(rustix::io::Errno::AGAIN) => {}
729            Err(_) => {
730                let _ = tx.send(None);
731                return;
732            }
733        }
734    }
735}
736
737/// Run to completion and capture output.
738pub(crate) fn run_captured(
739    client: &Client,
740    instance: &str,
741    argv: &[String],
742    base: &Request,
743    stdin: Stdin,
744    timeout: Option<Duration>,
745) -> Result<ExecOutput> {
746    let req = Request {
747        argv: argv.to_vec(),
748        tty: false,
749        ..base.clone()
750    };
751    start_with_timeout(client, instance, req, stdin, timeout)?.collect_output()
752}
753
754/// Run attached to this process's terminal: stdio forwarded, raw mode and window
755/// size in tty mode, signals forwarded. Returns the exit code.
756pub(crate) fn attach(
757    client: &Client,
758    instance: &str,
759    req: Request,
760    stdin: Stdin,
761    timeout: Option<Duration>,
762) -> Result<i32> {
763    use signal_hook::consts::signal::*;
764    let tty = req.tty;
765    let _raw = if tty { RawMode::enable() } else { None };
766    let mut stream = start_with_timeout(client, instance, req, stdin, timeout)?;
767    if tty {
768        if let Some((w, h)) = terminal_size() {
769            let _ = stream.resize(w, h);
770        }
771    }
772    let mut signals = signal_hook::iterator::Signals::new([
773        SIGINT, SIGTERM, SIGHUP, SIGQUIT, SIGUSR1, SIGUSR2, SIGWINCH,
774    ])?;
775    let handle = signals.handle();
776    let control = stream.control.clone();
777    let sig_thread = std::thread::spawn(move || {
778        for sig in signals.forever() {
779            if sig == SIGWINCH {
780                if tty {
781                    if let Some((w, h)) = terminal_size() {
782                        let _ = send_control(
783                            &control,
784                            json!({"command": "window-resize", "args": {"width": w.to_string(), "height": h.to_string()}}),
785                        );
786                    }
787                }
788                continue;
789            }
790            let _ = send_control(
791                &control,
792                json!({"command": "signal", "signal": linux_signal(sig)}),
793            );
794        }
795    });
796    let stdout = std::io::stdout();
797    let stderr = std::io::stderr();
798    while let Some(ev) = stream.next_event() {
799        match ev {
800            ExecEvent::Stdout(b) => {
801                let mut o = stdout.lock();
802                let _ = o.write_all(&b);
803                let _ = o.flush();
804            }
805            ExecEvent::Stderr(b) => {
806                let mut e = stderr.lock();
807                let _ = e.write_all(&b);
808                let _ = e.flush();
809            }
810        }
811    }
812    let code = stream.finish();
813    handle.close();
814    let _ = sig_thread.join();
815    code
816}
817
818/// The guest is Linux whatever the host is, and the host's numbering can
819/// differ: SIGUSR1/SIGUSR2 are 30/31 on macOS, 10/12 on Linux. The others
820/// forwarded (INT, TERM, HUP, QUIT) agree everywhere.
821fn linux_signal(sig: i32) -> i32 {
822    use signal_hook::consts::signal::{SIGUSR1, SIGUSR2};
823    match sig {
824        SIGUSR1 => 10,
825        SIGUSR2 => 12,
826        s => s,
827    }
828}
829
830/// Width and height of the terminal on stdout, if it is one.
831pub fn terminal_size() -> Option<(u16, u16)> {
832    let ws = rustix::termios::tcgetwinsize(std::io::stdout()).ok()?;
833    if ws.ws_col == 0 || ws.ws_row == 0 {
834        return None;
835    }
836    Some((ws.ws_col, ws.ws_row))
837}
838
839/// Whether both stdin and stdout are terminals (the default for `-t`).
840pub fn stdio_is_tty() -> bool {
841    rustix::termios::isatty(std::io::stdin()) && rustix::termios::isatty(std::io::stdout())
842}
843
844/// Puts the local terminal in raw mode for the life of the guard.
845struct RawMode {
846    saved: rustix::termios::Termios,
847}
848
849impl RawMode {
850    fn enable() -> Option<RawMode> {
851        let stdin = std::io::stdin();
852        let saved = rustix::termios::tcgetattr(&stdin).ok()?;
853        let mut raw = saved.clone();
854        raw.make_raw();
855        rustix::termios::tcsetattr(&stdin, rustix::termios::OptionalActions::Now, &raw).ok()?;
856        Some(RawMode { saved })
857    }
858}
859
860impl Drop for RawMode {
861    fn drop(&mut self) {
862        let _ = rustix::termios::tcsetattr(
863            std::io::stdin(),
864            rustix::termios::OptionalActions::Now,
865            &self.saved,
866        );
867    }
868}
869
870#[cfg(test)]
871mod tests {
872    use super::*;
873
874    #[test]
875    fn passwd_lines() {
876        let u = parse_passwd("dev:x:1000:1000:Dev,,,:/home/dev:/bin/bash\n").unwrap();
877        assert_eq!(u.uid, 1000);
878        assert_eq!(u.gid, 1000);
879        assert_eq!(u.home.as_deref(), Some("/home/dev"));
880        assert_eq!(u.shell.as_deref(), Some("/bin/bash"));
881        assert!(parse_passwd("short:x:1").is_none());
882        assert!(parse_passwd("bad:x:a:b:c:d:e").is_none());
883    }
884}