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
22pub type SpawnedProcess = (ProcessHandle, broadcast::Receiver<OutputChunk>);
24
25pub(crate) trait ChildTerminator: Send + 'static {
26 fn terminate(&mut self);
27}
28
29pub 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 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 pub fn output_subscribe(&self) -> broadcast::Receiver<OutputChunk> {
68 self.output_tx.subscribe()
69 }
70
71 pub async fn wait_for_exit(&self) {
73 loop {
74 if self.has_exited() {
75 return;
76 }
77 let notified = self.exit_notify.notified();
79 if self.has_exited() {
80 return;
81 }
82 notified.await;
83 }
84 }
85
86 pub fn has_exited(&self) -> bool {
88 lock_or_recover(&self.exit_code).is_some()
89 }
90
91 pub fn exit_code(&self) -> Option<i32> {
93 *lock_or_recover(&self.exit_code)
94 }
95
96 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}