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