Skip to main content

ante_exec/
handle.rs

1use anyhow::{Result, anyhow};
2use std::sync::{Arc, Mutex as StdMutex};
3use tokio::sync::{Notify, broadcast, mpsc};
4
5use crate::lock_or_recover;
6
7pub(crate) const OUTPUT_CHANNEL_CAPACITY: usize = 256;
8pub(crate) const STDIN_CHANNEL_CAPACITY: usize = 64;
9
10#[derive(Clone, Copy, Debug, Eq, PartialEq)]
11pub enum Stream {
12    Stdout,
13    Stderr,
14}
15
16#[derive(Clone, Debug, Eq, PartialEq)]
17pub struct OutputChunk {
18    pub stream: Stream,
19    pub data: Vec<u8>,
20}
21
22/// Handle + a pre-subscribed receiver that won't miss early output.
23pub type SpawnedProcess = (ProcessHandle, broadcast::Receiver<OutputChunk>);
24
25pub(crate) trait ChildTerminator: Send + 'static {
26    fn terminate(&mut self);
27}
28
29/// Thin process lifecycle handle.
30pub struct ProcessHandle {
31    output_tx: broadcast::Sender<OutputChunk>,
32    writer: Option<mpsc::Sender<Vec<u8>>>,
33    exit_code: Arc<StdMutex<Option<i32>>>,
34    exit_notify: Arc<Notify>,
35    terminator: StdMutex<Option<Box<dyn ChildTerminator>>>,
36}
37
38impl ProcessHandle {
39    pub(crate) fn from_parts(
40        output_tx: broadcast::Sender<OutputChunk>,
41        writer: Option<mpsc::Sender<Vec<u8>>>,
42        exit_code: Arc<StdMutex<Option<i32>>>,
43        exit_notify: Arc<Notify>,
44        terminator: Box<dyn ChildTerminator>,
45    ) -> Self {
46        Self {
47            output_tx,
48            writer,
49            exit_code,
50            exit_notify,
51            terminator: StdMutex::new(Some(terminator)),
52        }
53    }
54
55    /// Write bytes to the child's stdin. Errors if stdin was not piped.
56    pub async fn write_stdin(&self, data: &[u8]) -> Result<()> {
57        let writer =
58            self.writer.as_ref().ok_or_else(|| anyhow!("stdin was not piped for this process"))?;
59        writer
60            .send(data.to_vec())
61            .await
62            .map_err(|_| anyhow!("stdin is closed for this process"))?;
63        Ok(())
64    }
65
66    /// Subscribe to the raw output stream.
67    pub fn output_subscribe(&self) -> broadcast::Receiver<OutputChunk> {
68        self.output_tx.subscribe()
69    }
70
71    /// Wait until the process exits.
72    pub async fn wait_for_exit(&self) {
73        loop {
74            if self.has_exited() {
75                return;
76            }
77            // Acquire the waiter before checking again to avoid missing a notify.
78            let notified = self.exit_notify.notified();
79            if self.has_exited() {
80                return;
81            }
82            notified.await;
83        }
84    }
85
86    /// True if the child has exited.
87    pub fn has_exited(&self) -> bool {
88        lock_or_recover(&self.exit_code).is_some()
89    }
90
91    /// Exit code, if exited.
92    pub fn exit_code(&self) -> Option<i32> {
93        *lock_or_recover(&self.exit_code)
94    }
95
96    /// Kill the child and all descendants. Idempotent.
97    pub fn terminate(&self) {
98        let mut guard = lock_or_recover(&self.terminator);
99        if let Some(mut terminator) = guard.take() {
100            terminator.terminate();
101        }
102    }
103}
104
105impl Drop for ProcessHandle {
106    fn drop(&mut self) {
107        if !self.has_exited() {
108            self.terminate();
109        }
110    }
111}