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 '{marker}%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 marker = HANDLE_MARKER,
257 );
258 let output = ssh_output(target, &script)
259 .map_err(ServerLaunchError::from)
260 .map_err(LaunchFailure::from_error)?;
261 if !output.status.success() {
262 return Err(failed_ssh_handle_delivery(
263 target,
264 &remote_handle,
265 spec.container,
266 ServerLaunchError::Exit {
267 operation: format!("SSH launch on {target:?}"),
268 status: output.status,
269 diagnostics: String::from_utf8_lossy(&output.stderr).trim().to_owned(),
270 },
271 ));
272 }
273 let text = match String::from_utf8(output.stdout) {
274 Ok(text) => text,
275 Err(error) => {
276 return Err(failed_ssh_handle_delivery(
277 target,
278 &remote_handle,
279 spec.container,
280 ServerLaunchError::NonUtf8Identity {
281 target: target.to_owned(),
282 source: error,
283 },
284 ));
285 }
286 };
287 parse_ssh_handle(
288 target,
289 remote_stdout,
290 remote_stderr,
291 &text,
292 spec.container.map(str::to_owned),
293 )
294 .map_err(|error| failed_ssh_handle_delivery(target, &remote_handle, spec.container, error))
295}
296
297fn materialize_ssh_launch_files(
298 target: &str,
299 launch_files: &[LaunchFilePlan],
300) -> Result<(), ServerLaunchError> {
301 for launch_file in launch_files {
302 let script = remote_launch_file_script(launch_file)?;
303 let output = ssh_output_with_input(target, &script, launch_file.text.as_bytes())?;
304 if !output.status.success() {
305 return Err(ServerLaunchError::Exit {
306 operation: format!(
307 "materializing launch file {} over SSH on {target:?}",
308 launch_file.resolved_path.display()
309 ),
310 status: output.status,
311 diagnostics: String::from_utf8_lossy(&output.stderr).trim().to_owned(),
312 });
313 }
314 }
315 Ok(())
316}
317
318pub(super) fn remote_launch_file_script(
319 launch_file: &LaunchFilePlan,
320) -> Result<String, ServerLaunchError> {
321 let target = &launch_file.resolved_path;
322 let parent = target
323 .parent()
324 .ok_or_else(|| ServerLaunchError::Preparation {
325 message: format!(
326 "launch file target {} has no parent directory",
327 target.display()
328 ),
329 })?;
330 Ok(format!(
331 "# 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",
332 parent = shell_quote_path(parent),
333 target = shell_quote_path(target),
334 digest = shell_quote(&launch_file.sha256),
335 ))
336}
337
338fn parse_ssh_handle(
339 target: &str,
340 remote_stdout: PathBuf,
341 remote_stderr: PathBuf,
342 output: &str,
343 container: Option<String>,
344) -> Result<SshProcessHandle, ServerLaunchError> {
345 let result = output
346 .lines()
347 .rev()
348 .find_map(|line| line.strip_prefix(HANDLE_MARKER))
349 .ok_or_else(|| ServerLaunchError::MissingProcessId {
350 target: target.to_owned(),
351 })?
352 .split_once('\t')
353 .ok_or_else(|| ServerLaunchError::MissingStartTime {
354 target: target.to_owned(),
355 })?;
356 let leader_pid =
357 result
358 .0
359 .parse::<u32>()
360 .map_err(|source| ServerLaunchError::InvalidProcessId {
361 target: target.to_owned(),
362 value: result.0.to_owned(),
363 source,
364 })?;
365 let leader_start_time_ticks =
366 result
367 .1
368 .parse::<u64>()
369 .map_err(|source| ServerLaunchError::InvalidStartTime {
370 target: target.to_owned(),
371 value: result.1.to_owned(),
372 source,
373 })?;
374 let handle = SshProcessHandle {
375 target: target.to_owned(),
376 leader_pid,
377 process_group: leader_pid,
378 leader_start_time_ticks,
379 stdout: remote_stdout,
380 stderr: remote_stderr,
381 container,
382 };
383 handle
384 .validate()
385 .map_err(|details| ServerLaunchError::InvalidIdentity {
386 target: target.to_owned(),
387 details,
388 })?;
389 Ok(handle)
390}
391
392fn failed_ssh_handle_delivery(
393 target: &str,
394 remote_handle: &Path,
395 container: Option<&str>,
396 error: ServerLaunchError,
397) -> LaunchFailure {
398 let cleanup_started = Instant::now();
399 let process_cleanup = cleanup_incomplete_ssh_launch(target, remote_handle);
405 let removal = container.map(|container| remove_server_container(Some(target), container));
406 let removal_confirmed = removal.as_ref().is_none_or(|removal| removal.confirmed);
407 let (process_confirmed, cleanup_note) = match &process_cleanup {
408 Ok(()) => (
409 true,
410 format!(
411 "cleaned the remote process using {}",
412 remote_handle.display()
413 ),
414 ),
415 Err(cleanup) => (
416 false,
417 format!(
418 "remote launch cleanup using {} was not verified: {cleanup}",
419 remote_handle.display()
420 ),
421 ),
422 };
423 let verified = removal_confirmed && process_confirmed;
424 let mut cleanup = if process_confirmed {
425 completed_cleanup(CleanupTrigger::StartupRollback, false, false, Vec::new())
426 } else {
427 cleanup_error(
428 CleanupTrigger::StartupRollback,
429 false,
430 Vec::new(),
431 process_cleanup
432 .as_ref()
433 .err()
434 .cloned()
435 .unwrap_or_else(|| "remote launch cleanup was not verified".to_owned()),
436 )
437 };
438 cleanup.elapsed_ms = duration_millis(cleanup_started.elapsed());
439 cleanup.remote_deadline_ms = Some(duration_millis(REMOTE_SERVER_CLEANUP_DEADLINE));
440 cleanup.container_removal = removal.clone();
441 cleanup.verified = verified;
442 if !verified && cleanup.error.is_none() {
443 cleanup.error = removal.as_ref().and_then(|removal| removal.error.clone());
444 }
445 LaunchFailure {
446 error,
447 ownership_unknown: !verified,
448 container_removal: removal.map(Box::new),
449 cleanup: Some(Box::new(cleanup)),
450 cleanup_note: Some(cleanup_note),
451 }
452}
453
454fn cleanup_incomplete_ssh_launch(target: &str, remote_handle: &Path) -> Result<(), String> {
455 let bound = OperationBound::finite(REMOTE_SERVER_CLEANUP_DEADLINE);
456 let alive = remote_group_alive_script("$pid");
457 let script = format!(
458 "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\"",
459 file = shell_quote_path(remote_handle),
460 term_limit = TERM_POLL_LIMIT,
461 kill_limit = KILL_POLL_LIMIT,
462 );
463 let output = run_cleanup_command(
464 &ssh_argv(target, &script),
465 SSH_ENV_REMOVE,
466 &bound,
467 "SSH launch cleanup",
468 )
469 .map_err(|error| error.to_string())?;
470 if output.status.success() {
471 Ok(())
472 } else {
473 Err(format!(
474 "SSH cleanup exited with {}: {}",
475 output.status,
476 String::from_utf8_lossy(&output.stderr).trim()
477 ))
478 }
479}
480
481fn render_env_command(command: &CommandPlan) -> Result<String, String> {
482 if command.argv.is_empty() {
483 return Err("resolved server command is empty".to_owned());
484 }
485 let mut parts = vec!["env".to_owned(), "-i".to_owned()];
486 parts.extend(
487 command
488 .env
489 .iter()
490 .map(|(name, value)| shell_quote(&format!("{name}={value}"))),
491 );
492 parts.extend(
500 command
501 .pass_env
502 .iter()
503 .map(|name| format!("{name}=\"${{{name}}}\"")),
504 );
505 parts.extend(command.argv.iter().map(|value| shell_quote(value)));
506 Ok(parts.join(" "))
507}
508
509impl ProcessLauncher for SystemProcessRuntime {
510 fn spawn(&self, spec: ProcessSpec<'_>) -> Result<ProcessHandle, LaunchFailure> {
511 match spec.launch {
512 LaunchPlan::Local => spawn_local(spec).map(ProcessHandle::Local),
513 LaunchPlan::Ssh { target } => spawn_ssh(target, spec).map(ProcessHandle::Ssh),
514 }
515 }
516}