rskit_process/persistent/
run.rs1use 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
25pub 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 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 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}