Skip to main content

rskit_process/persistent/
mod.rs

1//! Persistent subprocess lifecycle support.
2
3use 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/// Captured output retained while waiting for persistent readiness.
47#[derive(Debug, Clone)]
48pub struct PersistentStartup {
49    /// Captured stdout at the moment readiness completed.
50    pub stdout: String,
51    /// Captured stdout bytes at the moment readiness completed.
52    pub stdout_bytes: Vec<u8>,
53    /// Captured stderr at the moment readiness completed.
54    pub stderr: String,
55    /// Captured stderr bytes at the moment readiness completed.
56    pub stderr_bytes: Vec<u8>,
57    /// Whether stdout capture exceeded the configured limit before readiness.
58    pub stdout_truncated: bool,
59    /// Whether stderr capture exceeded the configured limit before readiness.
60    pub stderr_truncated: bool,
61    /// Time elapsed from spawn until readiness completed.
62    pub duration: Duration,
63}
64
65/// Result of starting a persistent process.
66#[derive(Debug)]
67pub struct PersistentRun {
68    /// Startup output captured while waiting for readiness.
69    pub startup: PersistentStartup,
70    /// Running persistent process handle.
71    pub process: PersistentProcess,
72}
73
74/// Start a persistent process and wait for its readiness policy.
75pub 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                // Reuse the normal lifecycle cleanup so partially-started processes
143                // are terminated the same way as explicit shutdown.
144                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        // Reuse the normal lifecycle cleanup so readiness failures do not leak the child.
174        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}