Skip to main content

rskit_process/persistent/
run.rs

1use std::{
2    process::{Child, Command as StdCommand, Stdio},
3    sync::{Arc, atomic::AtomicBool, mpsc},
4    time::Instant,
5};
6
7use tokio_util::sync::CancellationToken;
8
9use crate::{
10    AppError, AppResult, EnvPolicy, ErrorCode, InputPolicy, ProcessConfig, ProcessIo, ProcessSpec,
11    SignalPolicy, process_group::isolate,
12};
13
14use super::cancel::spawn_cancel_thread;
15use super::config::PersistentConfig;
16use super::config::PersistentReadiness::{Command as CommandReadiness, OutputContains, Started};
17use super::error::{PersistentStartErrorKind, persistent_start_error};
18use super::io::{new_capture, spawn_output_readers, spawn_stdin_writer, take_capture};
19use super::process::{PersistentProcess, SpawnedProcess, cleanup_spawned_child, new_process};
20use super::readiness::{
21    readiness_wait_error, run_readiness_command, validate_readiness, wait_for_readiness,
22};
23use super::types::{PersistentRun, PersistentStartup};
24
25/// Start a persistent process and wait for its readiness policy.
26pub fn start_persistent_with_cancel(
27    spec: &ProcessSpec,
28    process_config: &ProcessConfig,
29    persistent_config: &PersistentConfig,
30    cancel: CancellationToken,
31) -> AppResult<PersistentRun> {
32    if spec.program.as_os_str().is_empty() {
33        return Err(AppError::invalid_input("program", "must not be empty"));
34    }
35    if cancel.is_cancelled() {
36        return Err(AppError::cancelled("persistent process startup"));
37    }
38    validate_readiness(&persistent_config.readiness)?;
39
40    let start = Instant::now();
41    let input = process_input(process_config)?;
42    let mut child = spawn_child(spec, process_config, input)?;
43    let stdout = new_capture();
44    let stderr = new_capture();
45    let cancelled = Arc::new(AtomicBool::new(false));
46    let cancel_thread = match spawn_cancel_thread(
47        child.id(),
48        cancel.clone(),
49        Arc::clone(&cancelled),
50        process_config.signal,
51        persistent_config.shutdown_grace_period,
52    ) {
53        Ok(thread) => Some(thread),
54        Err(error) => {
55            let _ = cleanup_spawned_child(
56                &mut child,
57                process_config.signal,
58                persistent_config.shutdown_grace_period,
59            );
60            return Err(error);
61        }
62    };
63    let (ready_tx, ready_rx) = mpsc::channel();
64    let (stdout_thread, stderr_thread) =
65        spawn_output_readers(&mut child, &stdout, &stderr, &ready_tx, persistent_config);
66    let stdin_thread = spawn_stdin_writer(&mut child, predefined_stdin(input));
67
68    let mut spawned = SpawnedProcess {
69        child,
70        stdin_thread,
71        stdout_thread,
72        stderr_thread,
73        cancel_thread,
74        cancelled,
75        stdout,
76        stderr,
77        start,
78    };
79
80    match &persistent_config.readiness {
81        Started => {
82            let _ = ready_tx.send(());
83        }
84        CommandReadiness(command) => {
85            if let Err(error) = run_readiness_command(
86                command,
87                process_config,
88                persistent_config.readiness_timeout,
89                persistent_config.shutdown_grace_period,
90                cancel.clone(),
91            ) {
92                let mut process =
93                    persistent_process(spawned, process_config.signal, persistent_config);
94                // Reuse the normal lifecycle cleanup
95                // so partially-started processes are terminated the same way as explicit shutdown.
96                let _ = process.shutdown_inner();
97                return Err(error);
98            }
99            let _ = ready_tx.send(());
100        }
101        OutputContains(_) => {}
102    }
103    drop(ready_tx);
104
105    if let Err(error) = wait_for_readiness(
106        &ready_rx,
107        persistent_config.readiness_timeout,
108        &cancel,
109        &spawned.cancelled,
110    ) {
111        let readiness_error = readiness_wait_error(&mut spawned.child, error, &spawned.cancelled)?;
112        let mut process = persistent_process(spawned, process_config.signal, persistent_config);
113        // Reuse the normal lifecycle cleanup so readiness failures do not leak the child.
114        let _ = process.shutdown_inner();
115        return Err(readiness_error);
116    }
117
118    let stdout_startup = take_capture(&spawned.stdout);
119    let stderr_startup = take_capture(&spawned.stderr);
120    let process = persistent_process(spawned, process_config.signal, persistent_config);
121
122    Ok(PersistentRun {
123        startup: PersistentStartup {
124            stdout: String::from_utf8_lossy(&stdout_startup.bytes).into_owned(),
125            stdout_bytes: stdout_startup.bytes,
126            stderr: String::from_utf8_lossy(&stderr_startup.bytes).into_owned(),
127            stderr_bytes: stderr_startup.bytes,
128            stdout_truncated: stdout_startup.truncated,
129            stderr_truncated: stderr_startup.truncated,
130            duration: start.elapsed(),
131        },
132        process,
133    })
134}
135
136fn process_input(config: &ProcessConfig) -> AppResult<&InputPolicy> {
137    match &config.io {
138        ProcessIo::Captured(io)
139            if matches!(io.input, InputPolicy::Closed | InputPolicy::Bytes(_)) =>
140        {
141            Ok(&io.input)
142        }
143        ProcessIo::Captured(_) => Err(AppError::invalid_input(
144            "process.io.input",
145            "persistent processes support only closed stdin or predefined stdin bytes",
146        )),
147        ProcessIo::Inherited(_) => Err(AppError::invalid_input(
148            "process.io",
149            "persistent processes use PersistentOutput for output handling; inherited mode is not supported",
150        )),
151        ProcessIo::Observed(_) => Err(AppError::invalid_input(
152            "process.io",
153            "persistent processes use PersistentOutput for observation; observed mode is not supported",
154        )),
155        #[cfg(unix)]
156        ProcessIo::Pty(_) => Err(AppError::invalid_input(
157            "process.io",
158            "persistent processes use PersistentOutput for output handling; pty mode is not supported",
159        )),
160    }
161}
162
163fn predefined_stdin(input: &InputPolicy) -> Option<Vec<u8>> {
164    match input {
165        InputPolicy::Bytes(bytes) => Some(bytes.clone()),
166        InputPolicy::Closed | InputPolicy::Inherit => None,
167    }
168}
169
170fn spawn_child(
171    spec: &ProcessSpec,
172    config: &ProcessConfig,
173    input: &InputPolicy,
174) -> AppResult<Child> {
175    let mut cmd = StdCommand::new(&spec.program);
176    cmd.args(&spec.args)
177        .stdin(if matches!(input, InputPolicy::Bytes(_)) {
178            Stdio::piped()
179        } else {
180            Stdio::null()
181        })
182        .stdout(Stdio::piped())
183        .stderr(Stdio::piped());
184
185    if let Some(dir) = &spec.dir {
186        cmd.current_dir(dir);
187    }
188    if matches!(spec.env_policy, EnvPolicy::Empty) {
189        cmd.env_clear();
190    }
191    for (key, value) in &spec.env {
192        cmd.env(key, value);
193    }
194    if config.signal.create_process_group {
195        isolate(&mut cmd);
196    }
197
198    cmd.spawn().map_err(|error| {
199        persistent_start_error(
200            PersistentStartErrorKind::SpawnFailed,
201            ErrorCode::Internal,
202            format!("failed to spawn persistent process: {error}"),
203        )
204        .with_cause(error)
205    })
206}
207
208fn persistent_process(
209    spawned: SpawnedProcess,
210    signal: SignalPolicy,
211    config: &PersistentConfig,
212) -> PersistentProcess {
213    new_process(spawned, signal, config.shutdown_grace_period)
214}