1use std::{
4 process::{Child, Command as StdCommand, Stdio},
5 sync::{Arc, atomic::AtomicBool, mpsc},
6 time::{Duration, Instant},
7};
8
9use tokio_util::sync::CancellationToken;
10
11use crate::{
12 AppError, AppResult, EnvPolicy, ErrorCode, InputPolicy, ProcessConfig, ProcessIo, ProcessSpec,
13 SignalPolicy, process_group::isolate,
14};
15
16mod cancel;
17mod config;
18mod error;
19mod io;
20mod process;
21mod readiness;
22
23#[cfg(all(test, unix))]
24mod tests;
25
26pub use config::{
27 PersistentConfig, PersistentOutput, PersistentOutputObserver, PersistentOutputStream,
28 PersistentReadiness,
29};
30pub use error::{PersistentStartErrorKind, persistent_start_error_kind};
31pub use process::{PersistentProcess, ShutdownOutcome};
32
33use cancel::spawn_cancel_thread;
34use config::PersistentReadiness::{Command as CommandReadiness, OutputContains, Started};
35use error::persistent_start_error;
36use io::{
37 ReaderThread, StdinThread, new_capture, spawn_output_readers, spawn_stdin_writer, take_capture,
38};
39use process::{cleanup_spawned_child, new_process};
40use readiness::{
41 readiness_wait_error, run_readiness_command, validate_readiness, wait_for_readiness,
42};
43
44#[derive(Debug, Clone)]
46pub struct PersistentStartup {
47 pub stdout: String,
49 pub stdout_bytes: Vec<u8>,
51 pub stderr: String,
53 pub stderr_bytes: Vec<u8>,
55 pub stdout_truncated: bool,
57 pub stderr_truncated: bool,
59 pub duration: Duration,
61}
62
63#[derive(Debug)]
65pub struct PersistentRun {
66 pub startup: PersistentStartup,
68 pub process: PersistentProcess,
70}
71
72pub fn start_persistent_with_cancel(
74 spec: &ProcessSpec,
75 process_config: &ProcessConfig,
76 persistent_config: &PersistentConfig,
77 cancel: CancellationToken,
78) -> AppResult<PersistentRun> {
79 if spec.program.as_os_str().is_empty() {
80 return Err(AppError::invalid_input("program", "must not be empty"));
81 }
82 if cancel.is_cancelled() {
83 return Err(AppError::cancelled("persistent process startup"));
84 }
85 validate_readiness(&persistent_config.readiness)?;
86
87 let start = Instant::now();
88 let input = process_input(process_config)?;
89 let mut child = spawn_child(spec, process_config, input)?;
90 let stdout = new_capture();
91 let stderr = new_capture();
92 let cancelled = Arc::new(AtomicBool::new(false));
93 let cancel_thread = match spawn_cancel_thread(
94 child.id(),
95 cancel.clone(),
96 Arc::clone(&cancelled),
97 process_config.signal,
98 persistent_config.shutdown_grace_period,
99 ) {
100 Ok(thread) => Some(thread),
101 Err(error) => {
102 let _ = cleanup_spawned_child(
103 &mut child,
104 process_config.signal,
105 persistent_config.shutdown_grace_period,
106 );
107 return Err(error);
108 }
109 };
110 let (ready_tx, ready_rx) = mpsc::channel();
111 let (stdout_thread, stderr_thread) =
112 spawn_output_readers(&mut child, &stdout, &stderr, &ready_tx, persistent_config);
113 let stdin_thread = spawn_stdin_writer(&mut child, predefined_stdin(input));
114
115 match &persistent_config.readiness {
116 Started => {
117 let _ = ready_tx.send(());
118 }
119 CommandReadiness(command) => {
120 if let Err(error) = run_readiness_command(
121 command,
122 process_config,
123 persistent_config.readiness_timeout,
124 persistent_config.shutdown_grace_period,
125 cancel.clone(),
126 ) {
127 let mut process = persistent_process(
128 child,
129 stdin_thread,
130 stdout_thread,
131 stderr_thread,
132 cancel_thread,
133 cancelled,
134 stdout,
135 stderr,
136 start,
137 process_config.signal,
138 persistent_config,
139 );
140 let _ = process.shutdown_inner();
143 return Err(error);
144 }
145 let _ = ready_tx.send(());
146 }
147 OutputContains(_) => {}
148 }
149 drop(ready_tx);
150
151 if let Err(error) = wait_for_readiness(
152 &ready_rx,
153 persistent_config.readiness_timeout,
154 &cancel,
155 &cancelled,
156 ) {
157 let readiness_error = readiness_wait_error(&mut child, error, &cancelled)?;
158 let mut process = persistent_process(
159 child,
160 stdin_thread,
161 stdout_thread,
162 stderr_thread,
163 cancel_thread,
164 cancelled,
165 stdout,
166 stderr,
167 start,
168 process_config.signal,
169 persistent_config,
170 );
171 let _ = process.shutdown_inner();
173 return Err(readiness_error);
174 }
175
176 let stdout_startup = take_capture(&stdout);
177 let stderr_startup = take_capture(&stderr);
178 let process = persistent_process(
179 child,
180 stdin_thread,
181 stdout_thread,
182 stderr_thread,
183 cancel_thread,
184 cancelled,
185 stdout,
186 stderr,
187 start,
188 process_config.signal,
189 persistent_config,
190 );
191
192 Ok(PersistentRun {
193 startup: PersistentStartup {
194 stdout: String::from_utf8_lossy(&stdout_startup.bytes).into_owned(),
195 stdout_bytes: stdout_startup.bytes,
196 stderr: String::from_utf8_lossy(&stderr_startup.bytes).into_owned(),
197 stderr_bytes: stderr_startup.bytes,
198 stdout_truncated: stdout_startup.truncated,
199 stderr_truncated: stderr_startup.truncated,
200 duration: start.elapsed(),
201 },
202 process,
203 })
204}
205
206fn process_input(config: &ProcessConfig) -> AppResult<&InputPolicy> {
207 match &config.io {
208 ProcessIo::Captured(io)
209 if matches!(io.input, InputPolicy::Closed | InputPolicy::Bytes(_)) =>
210 {
211 Ok(&io.input)
212 }
213 ProcessIo::Captured(_) => Err(AppError::invalid_input(
214 "process.io.input",
215 "persistent processes support only closed stdin or predefined stdin bytes",
216 )),
217 ProcessIo::Inherited(_) => Err(AppError::invalid_input(
218 "process.io",
219 "persistent processes use PersistentOutput for output handling; inherited mode is not supported",
220 )),
221 ProcessIo::Observed(_) => Err(AppError::invalid_input(
222 "process.io",
223 "persistent processes use PersistentOutput for observation; observed mode is not supported",
224 )),
225 #[cfg(unix)]
226 ProcessIo::Pty(_) => Err(AppError::invalid_input(
227 "process.io",
228 "persistent processes use PersistentOutput for output handling; pty mode is not supported",
229 )),
230 }
231}
232
233fn predefined_stdin(input: &InputPolicy) -> Option<Vec<u8>> {
234 match input {
235 InputPolicy::Bytes(bytes) => Some(bytes.clone()),
236 InputPolicy::Closed | InputPolicy::Inherit => None,
237 }
238}
239
240fn spawn_child(
241 spec: &ProcessSpec,
242 config: &ProcessConfig,
243 input: &InputPolicy,
244) -> AppResult<Child> {
245 let mut cmd = StdCommand::new(&spec.program);
246 cmd.args(&spec.args)
247 .stdin(if matches!(input, InputPolicy::Bytes(_)) {
248 Stdio::piped()
249 } else {
250 Stdio::null()
251 })
252 .stdout(Stdio::piped())
253 .stderr(Stdio::piped());
254
255 if let Some(dir) = &spec.dir {
256 cmd.current_dir(dir);
257 }
258 if matches!(spec.env_policy, EnvPolicy::Empty) {
259 cmd.env_clear();
260 }
261 for (key, value) in &spec.env {
262 cmd.env(key, value);
263 }
264 if config.signal.create_process_group {
265 isolate(&mut cmd);
266 }
267
268 cmd.spawn().map_err(|error| {
269 persistent_start_error(
270 PersistentStartErrorKind::SpawnFailed,
271 ErrorCode::Internal,
272 format!("failed to spawn persistent process: {error}"),
273 )
274 .with_cause(error)
275 })
276}
277
278#[allow(clippy::too_many_arguments)]
279fn persistent_process(
280 child: Child,
281 stdin_thread: StdinThread,
282 stdout_thread: ReaderThread,
283 stderr_thread: ReaderThread,
284 cancel_thread: Option<cancel::CancelThread>,
285 cancelled: Arc<AtomicBool>,
286 stdout: io::Capture,
287 stderr: io::Capture,
288 start: Instant,
289 signal: SignalPolicy,
290 config: &PersistentConfig,
291) -> PersistentProcess {
292 new_process(
293 child,
294 stdin_thread,
295 stdout_thread,
296 stderr_thread,
297 cancel_thread,
298 cancelled,
299 stdout,
300 stderr,
301 start,
302 signal,
303 config.shutdown_grace_period,
304 )
305}