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