1use super::cleanup::{
2 KILL_POLL_LIMIT, REMOTE_SERVER_CLEANUP_DEADLINE, TERM_POLL_LIMIT, cleanup_error,
3 cleanup_failed_local_launch, completed_cleanup, remove_server_container,
4};
5use super::observation::{remote_group_alive_script, run_cleanup_command};
6use super::{
7 CleanupTrigger, HostProcessHandle, LaunchFailure, ProcessHandle, ProcessLauncher, ProcessSpec,
8 ServerLaunchError, SshProcessHandle, SystemProcessRuntime,
9};
10use crate::operation_bound::{OperationBound, duration_millis};
11use crate::plan::{CommandPlan, LaunchFilePlan, LaunchPlan};
12use crate::shell::{shell_quote, shell_quote_path};
13use crate::ssh::{SSH_ENV_REMOVE, ssh_argv, ssh_output, ssh_output_with_input};
14use sha2::{Digest, Sha256};
15use std::fs;
16use std::fs::File;
17use std::io::{self, Read, Write};
18use std::os::unix::fs::PermissionsExt;
19use std::os::unix::process::CommandExt;
20use std::path::{Path, PathBuf};
21use std::process::{Command, Stdio};
22use std::time::Instant;
23
24const HANDLE_MARKER: &str = "INFERLAB_HANDLE\t";
25
26pub(super) fn spawn_local(spec: ProcessSpec<'_>) -> Result<HostProcessHandle, LaunchFailure> {
27 let fail = |message: String| LaunchFailure::before_launch(message);
28 fs::create_dir_all(spec.cache_root).map_err(|source| {
29 LaunchFailure::from_error(ServerLaunchError::FileIo {
30 operation: "create runtime cache root",
31 path: spec.cache_root.to_path_buf(),
32 source,
33 })
34 })?;
35 materialize_local_launch_files(spec.launch_files).map_err(LaunchFailure::from_error)?;
36 let (program, args) = spec
37 .command
38 .argv
39 .split_first()
40 .ok_or_else(|| fail("resolved server command is empty".to_owned()))?;
41 let stdout = File::create(spec.stdout).map_err(|source| {
42 LaunchFailure::from_error(ServerLaunchError::FileIo {
43 operation: "create server stdout log",
44 path: spec.stdout.to_path_buf(),
45 source,
46 })
47 })?;
48 let stderr = File::create(spec.stderr).map_err(|source| {
49 LaunchFailure::from_error(ServerLaunchError::FileIo {
50 operation: "create server stderr log",
51 path: spec.stderr.to_path_buf(),
52 source,
53 })
54 })?;
55 let mut command = Command::new(program);
56 command
57 .args(args)
58 .current_dir(&spec.command.cwd)
59 .env_clear()
60 .envs(&spec.command.env)
61 .stdin(Stdio::null())
62 .stdout(Stdio::from(stdout))
63 .stderr(Stdio::from(stderr))
64 .process_group(0);
65 for name in &spec.command.pass_env {
73 if let Some(value) = std::env::var_os(name) {
74 command.env(name, value);
75 }
76 }
77 let mut child = command.spawn().map_err(|source| {
78 LaunchFailure::from_error(ServerLaunchError::Process {
79 program: program.clone(),
80 source,
81 })
82 })?;
83 HostProcessHandle::new(child.id(), spec.container.map(str::to_owned)).map_err(|error| {
84 let mut cleanup = cleanup_failed_local_launch(&mut child);
85 let error = ServerLaunchError::Preparation { message: error };
86 match spec.container {
92 Some(container) => {
93 let removal = remove_server_container(None, container);
94 if !removal.confirmed {
95 cleanup.verified = false;
96 if cleanup.error.is_none() {
97 cleanup.error = removal.error.clone();
98 }
99 }
100 cleanup.container_removal = Some(removal.clone());
101 LaunchFailure {
102 error,
103 ownership_unknown: !cleanup.verified,
104 container_removal: Some(Box::new(removal)),
105 cleanup: Some(Box::new(cleanup)),
106 cleanup_note: None,
107 }
108 }
109 None => LaunchFailure {
110 error,
111 ownership_unknown: !cleanup.verified,
112 container_removal: None,
113 cleanup: Some(Box::new(cleanup)),
114 cleanup_note: None,
115 },
116 }
117 })
118}
119
120pub(super) fn materialize_local_launch_files(
121 launch_files: &[LaunchFilePlan],
122) -> Result<(), ServerLaunchError> {
123 for launch_file in launch_files {
124 publish_local_launch_file(launch_file)?;
125 }
126 Ok(())
127}
128
129fn publish_local_launch_file(launch_file: &LaunchFilePlan) -> Result<(), ServerLaunchError> {
130 let target = &launch_file.resolved_path;
131 let parent = target
132 .parent()
133 .ok_or_else(|| ServerLaunchError::Preparation {
134 message: format!(
135 "launch file target {} has no parent directory",
136 target.display()
137 ),
138 })?;
139 fs::create_dir_all(parent).map_err(|source| ServerLaunchError::FileIo {
140 operation: "create launch file directory",
141 path: parent.to_path_buf(),
142 source,
143 })?;
144 let mut staged = tempfile::Builder::new()
145 .prefix(".inferlab-launch.")
146 .tempfile_in(parent)
147 .map_err(|source| ServerLaunchError::FileIo {
148 operation: "stage launch file",
149 path: target.to_path_buf(),
150 source,
151 })?;
152 staged
153 .write_all(launch_file.text.as_bytes())
154 .and_then(|()| staged.flush())
155 .map_err(|source| ServerLaunchError::FileIo {
156 operation: "write staged launch file",
157 path: target.to_path_buf(),
158 source,
159 })?;
160 staged
161 .as_file()
162 .set_permissions(fs::Permissions::from_mode(0o444))
163 .and_then(|()| staged.as_file().sync_all())
164 .map_err(|source| ServerLaunchError::FileIo {
165 operation: "finalize staged launch file",
166 path: target.to_path_buf(),
167 source,
168 })?;
169
170 match staged.persist_noclobber(target) {
171 Ok(_) => Ok(()),
172 Err(failure) => {
173 let source = failure.error;
174 if source.kind() == io::ErrorKind::AlreadyExists {
175 verify_existing_launch_file(target, &launch_file.sha256)
176 } else {
177 Err(ServerLaunchError::FileIo {
178 operation: "publish launch file",
179 path: target.to_path_buf(),
180 source,
181 })
182 }
183 }
184 }
185}
186
187fn verify_existing_launch_file(
188 target: &Path,
189 expected_sha256: &str,
190) -> Result<(), ServerLaunchError> {
191 let metadata = fs::symlink_metadata(target).map_err(|source| ServerLaunchError::FileIo {
192 operation: "inspect existing launch file",
193 path: target.to_path_buf(),
194 source,
195 })?;
196 if !metadata.file_type().is_file() {
197 return Err(ServerLaunchError::NotRegularFile {
198 path: target.to_path_buf(),
199 });
200 }
201 let actual_sha256 = file_sha256(target)?;
202 if actual_sha256 != expected_sha256 {
203 return Err(ServerLaunchError::FileDigestMismatch {
204 path: target.to_path_buf(),
205 expected: expected_sha256.to_owned(),
206 actual: actual_sha256,
207 });
208 }
209 Ok(())
210}
211
212fn file_sha256(path: &Path) -> Result<String, ServerLaunchError> {
213 let mut file = File::open(path).map_err(|source| ServerLaunchError::FileIo {
214 operation: "read launch file",
215 path: path.to_path_buf(),
216 source,
217 })?;
218 let mut digest = Sha256::new();
219 let mut buffer = [0_u8; 8192];
220 loop {
221 let read = file
222 .read(&mut buffer)
223 .map_err(|source| ServerLaunchError::FileIo {
224 operation: "read launch file",
225 path: path.to_path_buf(),
226 source,
227 })?;
228 if read == 0 {
229 break;
230 }
231 digest.update(&buffer[..read]);
232 }
233 Ok(format!("{:x}", digest.finalize()))
234}
235
236pub(super) fn spawn_ssh(
240 target: &str,
241 spec: ProcessSpec<'_>,
242) -> Result<SshProcessHandle, LaunchFailure> {
243 let remote_stdout = spec.remote_dir.join("stdout.log");
244 let remote_stderr = spec.remote_dir.join("stderr.log");
245 let remote_handle = spec.remote_dir.join("launch.handle");
246 let command = render_env_command(spec.command).map_err(LaunchFailure::before_launch)?;
247 materialize_ssh_launch_files(target, spec.launch_files).map_err(LaunchFailure::from_error)?;
248 let script = format!(
249 "set -eu; mkdir -p {dir} {cache}; cd {cwd}; nohup setsid {command} >{stdout} 2>{stderr} </dev/null & pid=$!; cleanup_pending=1; cleanup_launch() {{ if [ \"$cleanup_pending\" = 1 ]; then kill -KILL -- -$pid 2>/dev/null || kill -KILL $pid 2>/dev/null || true; fi; }}; trap cleanup_launch EXIT; ticks=$(awk '{{print $22}}' /proc/$pid/stat); printf '%s %s\\n' \"$pid\" \"$ticks\" > {handle}; printf 'INFERLAB_HANDLE\\t%s\\t%s\\n' \"$pid\" \"$ticks\"; cleanup_pending=0; trap - EXIT",
250 dir = shell_quote_path(spec.remote_dir),
251 cache = shell_quote_path(spec.cache_root),
252 cwd = shell_quote_path(&spec.command.cwd),
253 stdout = shell_quote_path(&remote_stdout),
254 stderr = shell_quote_path(&remote_stderr),
255 handle = shell_quote_path(&remote_handle),
256 );
257 let output = ssh_output(target, &script)
258 .map_err(ServerLaunchError::from)
259 .map_err(LaunchFailure::from_error)?;
260 if !output.status.success() {
261 return Err(failed_ssh_handle_delivery(
262 target,
263 &remote_handle,
264 spec.container,
265 ServerLaunchError::Exit {
266 operation: format!("SSH launch on {target:?}"),
267 status: output.status,
268 diagnostics: String::from_utf8_lossy(&output.stderr).trim().to_owned(),
269 },
270 ));
271 }
272 let text = match String::from_utf8(output.stdout) {
273 Ok(text) => text,
274 Err(error) => {
275 return Err(failed_ssh_handle_delivery(
276 target,
277 &remote_handle,
278 spec.container,
279 ServerLaunchError::NonUtf8Identity {
280 target: target.to_owned(),
281 source: error,
282 },
283 ));
284 }
285 };
286 parse_ssh_handle(
287 target,
288 remote_stdout,
289 remote_stderr,
290 &text,
291 spec.container.map(str::to_owned),
292 )
293 .map_err(|error| failed_ssh_handle_delivery(target, &remote_handle, spec.container, error))
294}
295
296fn materialize_ssh_launch_files(
297 target: &str,
298 launch_files: &[LaunchFilePlan],
299) -> Result<(), ServerLaunchError> {
300 for launch_file in launch_files {
301 let script = remote_launch_file_script(launch_file)?;
302 let output = ssh_output_with_input(target, &script, launch_file.text.as_bytes())?;
303 if !output.status.success() {
304 return Err(ServerLaunchError::Exit {
305 operation: format!(
306 "materializing launch file {} over SSH on {target:?}",
307 launch_file.resolved_path.display()
308 ),
309 status: output.status,
310 diagnostics: String::from_utf8_lossy(&output.stderr).trim().to_owned(),
311 });
312 }
313 }
314 Ok(())
315}
316
317pub(super) fn remote_launch_file_script(
318 launch_file: &LaunchFilePlan,
319) -> Result<String, ServerLaunchError> {
320 let target = &launch_file.resolved_path;
321 let parent = target
322 .parent()
323 .ok_or_else(|| ServerLaunchError::Preparation {
324 message: format!(
325 "launch file target {} has no parent directory",
326 target.display()
327 ),
328 })?;
329 Ok(format!(
330 "# INFERLAB_LAUNCH_FILE\nset -eu\numask 077\nparent={parent}\ntarget={target}\ndigest={digest}\nmkdir -p -- \"$parent\"\nstage=$(mktemp \"$parent/.inferlab-launch.XXXXXX\")\ntrap 'rm -f -- \"$stage\"' EXIT\ncat > \"$stage\"\nactual=$(sha256sum -- \"$stage\" | awk '{{print $1}}')\nif [ \"$actual\" != \"$digest\" ]; then printf 'staged launch file digest mismatch for %s: expected %s, found %s\\n' \"$target\" \"$digest\" \"$actual\" >&2; exit 1; fi\nchmod 0444 -- \"$stage\"\nif ln -T -- \"$stage\" \"$target\" 2>/dev/null; then exit 0; fi\nif [ ! -f \"$target\" ] || [ -L \"$target\" ]; then printf 'existing launch file target %s is not a regular file\\n' \"$target\" >&2; exit 1; fi\nactual=$(sha256sum -- \"$target\" | awk '{{print $1}}')\nif [ \"$actual\" != \"$digest\" ]; then printf 'existing launch file %s does not match declared digest %s; found %s\\n' \"$target\" \"$digest\" \"$actual\" >&2; exit 1; fi",
331 parent = shell_quote_path(parent),
332 target = shell_quote_path(target),
333 digest = shell_quote(&launch_file.sha256),
334 ))
335}
336
337fn parse_ssh_handle(
338 target: &str,
339 remote_stdout: PathBuf,
340 remote_stderr: PathBuf,
341 output: &str,
342 container: Option<String>,
343) -> Result<SshProcessHandle, ServerLaunchError> {
344 let result = output
345 .lines()
346 .rev()
347 .find_map(|line| line.strip_prefix(HANDLE_MARKER))
348 .ok_or_else(|| ServerLaunchError::MissingProcessId {
349 target: target.to_owned(),
350 })?
351 .split_once('\t')
352 .ok_or_else(|| ServerLaunchError::MissingStartTime {
353 target: target.to_owned(),
354 })?;
355 let leader_pid =
356 result
357 .0
358 .parse::<u32>()
359 .map_err(|source| ServerLaunchError::InvalidProcessId {
360 target: target.to_owned(),
361 value: result.0.to_owned(),
362 source,
363 })?;
364 let leader_start_time_ticks =
365 result
366 .1
367 .parse::<u64>()
368 .map_err(|source| ServerLaunchError::InvalidStartTime {
369 target: target.to_owned(),
370 value: result.1.to_owned(),
371 source,
372 })?;
373 let handle = SshProcessHandle {
374 target: target.to_owned(),
375 leader_pid,
376 process_group: leader_pid,
377 leader_start_time_ticks,
378 stdout: remote_stdout,
379 stderr: remote_stderr,
380 container,
381 };
382 handle
383 .validate()
384 .map_err(|details| ServerLaunchError::InvalidIdentity {
385 target: target.to_owned(),
386 details,
387 })?;
388 Ok(handle)
389}
390
391fn failed_ssh_handle_delivery(
392 target: &str,
393 remote_handle: &Path,
394 container: Option<&str>,
395 error: ServerLaunchError,
396) -> LaunchFailure {
397 let cleanup_started = Instant::now();
398 let process_cleanup = cleanup_incomplete_ssh_launch(target, remote_handle);
404 let removal = container.map(|container| remove_server_container(Some(target), container));
405 let removal_confirmed = removal.as_ref().is_none_or(|removal| removal.confirmed);
406 let (process_confirmed, cleanup_note) = match &process_cleanup {
407 Ok(()) => (
408 true,
409 format!(
410 "cleaned the remote process using {}",
411 remote_handle.display()
412 ),
413 ),
414 Err(cleanup) => (
415 false,
416 format!(
417 "remote launch cleanup using {} was not verified: {cleanup}",
418 remote_handle.display()
419 ),
420 ),
421 };
422 let verified = removal_confirmed && process_confirmed;
423 let mut cleanup = if process_confirmed {
424 completed_cleanup(CleanupTrigger::StartupRollback, false, false, Vec::new())
425 } else {
426 cleanup_error(
427 CleanupTrigger::StartupRollback,
428 false,
429 Vec::new(),
430 process_cleanup
431 .as_ref()
432 .err()
433 .cloned()
434 .unwrap_or_else(|| "remote launch cleanup was not verified".to_owned()),
435 )
436 };
437 cleanup.elapsed_ms = duration_millis(cleanup_started.elapsed());
438 cleanup.remote_deadline_ms = Some(duration_millis(REMOTE_SERVER_CLEANUP_DEADLINE));
439 cleanup.container_removal = removal.clone();
440 cleanup.verified = verified;
441 if !verified && cleanup.error.is_none() {
442 cleanup.error = removal.as_ref().and_then(|removal| removal.error.clone());
443 }
444 LaunchFailure {
445 error,
446 ownership_unknown: !verified,
447 container_removal: removal.map(Box::new),
448 cleanup: Some(Box::new(cleanup)),
449 cleanup_note: Some(cleanup_note),
450 }
451}
452
453fn cleanup_incomplete_ssh_launch(target: &str, remote_handle: &Path) -> Result<(), String> {
454 let bound = OperationBound::finite(REMOTE_SERVER_CLEANUP_DEADLINE);
455 let alive = remote_group_alive_script("$pid");
456 let script = format!(
457 "set +e; file={file}; if [ ! -r \"$file\" ]; then exit 4; fi; read pid expected < \"$file\" || exit 4; if [ -r /proc/$pid/stat ]; then actual=$(awk '{{print $22}}' /proc/$pid/stat) || exit 4; [ \"$actual\" = \"$expected\" ] || exit 4; elif {alive}; then exit 5; else rm -f \"$file\"; exit 0; fi; if ! {alive}; then rm -f \"$file\"; exit 0; fi; kill -TERM -- -$pid; i=0; while {alive} && [ $i -lt {term_limit} ]; do sleep 0.1; i=$((i+1)); done; if {alive}; then kill -KILL -- -$pid; i=0; while {alive} && [ $i -lt {kill_limit} ]; do sleep 0.1; i=$((i+1)); done; fi; if {alive}; then exit 6; fi; rm -f \"$file\"",
458 file = shell_quote_path(remote_handle),
459 term_limit = TERM_POLL_LIMIT,
460 kill_limit = KILL_POLL_LIMIT,
461 );
462 let output = run_cleanup_command(
463 &ssh_argv(target, &script),
464 SSH_ENV_REMOVE,
465 &bound,
466 "SSH launch cleanup",
467 )
468 .map_err(|error| error.to_string())?;
469 if output.status.success() {
470 Ok(())
471 } else {
472 Err(format!(
473 "SSH cleanup exited with {}: {}",
474 output.status,
475 String::from_utf8_lossy(&output.stderr).trim()
476 ))
477 }
478}
479
480fn render_env_command(command: &CommandPlan) -> Result<String, String> {
481 if command.argv.is_empty() {
482 return Err("resolved server command is empty".to_owned());
483 }
484 let mut parts = vec!["env".to_owned(), "-i".to_owned()];
485 parts.extend(
486 command
487 .env
488 .iter()
489 .map(|(name, value)| shell_quote(&format!("{name}={value}"))),
490 );
491 parts.extend(
499 command
500 .pass_env
501 .iter()
502 .map(|name| format!("{name}=\"${{{name}}}\"")),
503 );
504 parts.extend(command.argv.iter().map(|value| shell_quote(value)));
505 Ok(parts.join(" "))
506}
507
508impl ProcessLauncher for SystemProcessRuntime {
509 fn spawn(&self, spec: ProcessSpec<'_>) -> Result<ProcessHandle, LaunchFailure> {
510 match spec.launch {
511 LaunchPlan::Local => spawn_local(spec).map(ProcessHandle::Local),
512 LaunchPlan::Ssh { target } => spawn_ssh(target, spec).map(ProcessHandle::Ssh),
513 }
514 }
515}