Skip to main content

inferlab_runtime/server/
observation.rs

1use super::{
2    HostProcessHandle, LogSyncError, ProcessCommandError, ProcessHandle, ProcessObserver,
3    ProcessStatus, REMOTE_LOG_SYNC_DEADLINE, SshProcessHandle, SystemProcessRuntime,
4};
5use crate::operation_bound::OperationBound;
6use crate::process_group::process_start_time;
7use crate::shell::shell_quote_path;
8use crate::ssh::{SSH_ENV_REMOVE, ssh_argv, ssh_output};
9use std::fs;
10use std::path::Path;
11use std::process::{Command, Output};
12
13pub(super) fn verified_local_status(handle: &HostProcessHandle) -> ProcessStatus {
14    verified_local_status_under(handle, None, false)
15}
16
17pub(super) fn verified_local_status_with_bound(
18    handle: &HostProcessHandle,
19    bound: &OperationBound,
20) -> ProcessStatus {
21    verified_local_status_under(handle, Some(bound), false)
22}
23
24pub(super) fn verified_local_status_under(
25    handle: &HostProcessHandle,
26    bound: Option<&OperationBound>,
27    cleanup: bool,
28) -> ProcessStatus {
29    if let Err(error) = handle.validate() {
30        return status_error(error);
31    }
32    match process_start_time(handle.leader_pid) {
33        Ok(Some(actual)) if actual != handle.leader_start_time_ticks => ProcessStatus {
34            queried: true,
35            alive: false,
36            error: Some(format!(
37                "managed process {} exited and its pid was reused: recorded start time {}, observed {}",
38                handle.leader_pid, handle.leader_start_time_ticks, actual
39            )),
40        },
41        Ok(Some(_)) => {
42            match process_group_has_live_members_under(handle.process_group, bound, cleanup) {
43                Ok(alive) => ProcessStatus {
44                    queried: true,
45                    alive,
46                    error: None,
47                },
48                Err(error) => status_error(error.to_string()),
49            }
50        }
51        Ok(None) => {
52            match process_group_has_live_members_under(handle.process_group, bound, cleanup) {
53                Ok(false) => ProcessStatus {
54                    queried: true,
55                    alive: false,
56                    error: None,
57                },
58                Ok(true) => status_error(format!(
59                    "process-group {} still has members but recorded leader {} no longer exists; ownership cannot be verified",
60                    handle.process_group, handle.leader_pid
61                )),
62                Err(error) => status_error(error.to_string()),
63            }
64        }
65        Err(error) => status_error(error.to_string()),
66    }
67}
68
69pub(super) fn verified_ssh_status(handle: &SshProcessHandle) -> ProcessStatus {
70    verified_ssh_status_under(handle, None)
71}
72
73pub(super) fn verified_ssh_status_with_bound(
74    handle: &SshProcessHandle,
75    bound: &OperationBound,
76) -> ProcessStatus {
77    verified_ssh_status_under(handle, Some(bound))
78}
79
80pub(super) fn verified_ssh_status_under(
81    handle: &SshProcessHandle,
82    bound: Option<&OperationBound>,
83) -> ProcessStatus {
84    if let Err(error) = handle.validate() {
85        return status_error(error);
86    }
87    let script = format!(
88        "set -eu; pid={}; expected={}; if [ -r /proc/$pid/stat ]; then actual=$(awk '{{print $22}}' /proc/$pid/stat); if [ \"$actual\" != \"$expected\" ]; then printf 'stale %s\\n' \"$actual\"; exit 4; fi; elif {}; then printf 'unknown leader-missing\\n'; exit 5; else printf 'dead\\n'; exit 3; fi; if {}; then printf 'alive\\n'; exit 0; fi; printf 'dead\\n'; exit 3",
89        handle.leader_pid,
90        handle.leader_start_time_ticks,
91        remote_group_alive_script(&handle.process_group.to_string()),
92        remote_group_alive_script(&handle.process_group.to_string()),
93    );
94    let output = match bound {
95        Some(bound) => {
96            run_status_command(&ssh_argv(&handle.target, &script), SSH_ENV_REMOVE, bound)
97        }
98        None => ssh_output(&handle.target, &script).map_err(|source| ProcessCommandError::Ssh {
99            operation: "process status command".to_owned(),
100            source,
101        }),
102    };
103    match output {
104        Ok(output) if output.status.success() => ProcessStatus {
105            queried: true,
106            alive: true,
107            error: None,
108        },
109        Ok(output) if output.status.code() == Some(3) => ProcessStatus {
110            queried: true,
111            alive: false,
112            error: None,
113        },
114        Ok(output) if output.status.code() == Some(4) => ProcessStatus {
115            queried: true,
116            alive: false,
117            error: Some(format!(
118                "managed SSH process {} exited and its pid was reused: {}",
119                handle.leader_pid,
120                String::from_utf8_lossy(&output.stdout).trim()
121            )),
122        },
123        Ok(output) if output.status.code() == Some(5) => status_error(format!(
124            "SSH process-group {} ownership could not be verified: {}",
125            handle.process_group,
126            String::from_utf8_lossy(&output.stdout).trim()
127        )),
128        Ok(output) => status_error(format!(
129            "SSH status exited with {}: {}",
130            output.status,
131            String::from_utf8_lossy(&output.stderr).trim()
132        )),
133        Err(error) => status_error(error.to_string()),
134    }
135}
136
137fn status_error(error: String) -> ProcessStatus {
138    ProcessStatus {
139        queried: false,
140        alive: false,
141        error: Some(error),
142    }
143}
144
145pub(super) fn remote_group_alive_script(group: &str) -> String {
146    format!(
147        "ps -eo pgid=,stat= | awk -v pgid={group} '$1 == pgid && $2 !~ /^Z/ {{ found=1 }} END {{ exit !found }}'"
148    )
149}
150
151pub(super) fn fetch_remote_file(
152    target: &str,
153    remote: &Path,
154    local: &Path,
155    bound: &OperationBound,
156    cleanup: bool,
157) -> Result<(), LogSyncError> {
158    let argv = ssh_argv(target, &format!("cat -- {}", shell_quote_path(remote)));
159    let output = if cleanup {
160        run_cleanup_command(&argv, SSH_ENV_REMOVE, bound, "remote log synchronization")
161    } else {
162        run_status_command(&argv, SSH_ENV_REMOVE, bound)
163    }
164    .map_err(|source| LogSyncError::ReadRemote {
165        path: remote.to_path_buf(),
166        source,
167    })?;
168    if !output.status.success() {
169        return Err(LogSyncError::RemoteExit {
170            path: remote.to_path_buf(),
171            status: output.status,
172            stderr: String::from_utf8_lossy(&output.stderr).trim().to_owned(),
173        });
174    }
175    fs::write(local, output.stdout).map_err(|source| LogSyncError::WriteLocal {
176        path: local.to_path_buf(),
177        source,
178    })
179}
180
181fn process_group_has_live_members_under(
182    process_group: u32,
183    bound: Option<&OperationBound>,
184    cleanup: bool,
185) -> Result<bool, ProcessCommandError> {
186    let argv = ["ps", "-eo", "pid=,pgid=,stat="];
187    let output = match bound {
188        Some(bound) if cleanup => run_cleanup_command(&argv, &[], bound, "process cleanup status"),
189        Some(bound) => run_status_command(&argv, &[], bound),
190        None => Command::new(argv[0])
191            .args(&argv[1..])
192            .output()
193            .map_err(|source| ProcessCommandError::Launch {
194                operation: "process-group query".to_owned(),
195                source,
196            }),
197    }?;
198    if !output.status.success() {
199        return Err(ProcessCommandError::Exit {
200            operation: "process-group query".to_owned(),
201            status: output.status,
202            stderr: String::from_utf8_lossy(&output.stderr).trim().to_owned(),
203        });
204    }
205    let process_group = process_group.to_string();
206    Ok(String::from_utf8_lossy(&output.stdout)
207        .lines()
208        .filter_map(|line| {
209            let mut fields = line.split_whitespace();
210            let _pid = fields.next()?;
211            let group = fields.next()?;
212            let state = fields.next()?;
213            Some((group, state))
214        })
215        .any(|(group, state)| group == process_group && !state.starts_with('Z')))
216}
217
218pub(super) fn run_status_command<S: AsRef<std::ffi::OsStr>>(
219    argv: &[S],
220    env_remove: &[&str],
221    bound: &OperationBound,
222) -> Result<Output, ProcessCommandError> {
223    let operation = "process status command";
224    match crate::container::run_with_bound(argv, env_remove, None, None, bound, None) {
225        Ok(crate::container::BoundedWait::Exited {
226            status,
227            stdout,
228            stderr,
229        }) => Ok(Output {
230            status,
231            stdout,
232            stderr,
233        }),
234        Ok(crate::container::BoundedWait::Expired { kill, .. }) => {
235            kill.map_err(|source| ProcessCommandError::Io {
236                operation: "process status cleanup".to_owned(),
237                source,
238            })?;
239            Err(ProcessCommandError::Deadline {
240                operation: "process status attempt".to_owned(),
241            })
242        }
243        Ok(crate::container::BoundedWait::Interrupted { kill, .. }) => {
244            kill.map_err(|source| ProcessCommandError::Io {
245                operation: "process status cleanup".to_owned(),
246                source,
247            })?;
248            Err(ProcessCommandError::Interrupted {
249                operation: "process status attempt".to_owned(),
250            })
251        }
252        Err(crate::container::BoundedError::Launch(source)) => Err(ProcessCommandError::Launch {
253            operation: operation.to_owned(),
254            source,
255        }),
256        Err(
257            crate::container::BoundedError::Stdin(error)
258            | crate::container::BoundedError::Wait(error),
259        ) => Err(ProcessCommandError::Io {
260            operation: operation.to_owned(),
261            source: error,
262        }),
263        Err(crate::container::BoundedError::WaitCleanup {
264            source, cleanup, ..
265        }) => Err(ProcessCommandError::WaitCleanup {
266            operation: operation.to_owned(),
267            source,
268            cleanup: cleanup.error.unwrap_or_else(|| {
269                if cleanup.verified {
270                    "verified"
271                } else {
272                    "unverified"
273                }
274                .to_owned()
275            }),
276        }),
277    }
278}
279
280pub(super) fn run_cleanup_command<S: AsRef<std::ffi::OsStr>>(
281    argv: &[S],
282    env_remove: &[&str],
283    bound: &OperationBound,
284    operation: &str,
285) -> Result<Output, ProcessCommandError> {
286    match crate::container::run_cleanup_with_bound(argv, env_remove, None, None, bound, None) {
287        Ok(crate::container::BoundedWait::Exited {
288            status,
289            stdout,
290            stderr,
291        }) => Ok(Output {
292            status,
293            stdout,
294            stderr,
295        }),
296        Ok(crate::container::BoundedWait::Expired { kill, .. }) => {
297            kill.map_err(|source| ProcessCommandError::Io {
298                operation: format!("{operation} child cleanup"),
299                source,
300            })?;
301            Err(ProcessCommandError::Deadline {
302                operation: operation.to_owned(),
303            })
304        }
305        Ok(crate::container::BoundedWait::Interrupted { kill, .. }) => {
306            kill.map_err(|source| ProcessCommandError::Io {
307                operation: format!("{operation} child cleanup"),
308                source,
309            })?;
310            Err(ProcessCommandError::Interrupted {
311                operation: operation.to_owned(),
312            })
313        }
314        Err(crate::container::BoundedError::Launch(source)) => Err(ProcessCommandError::Launch {
315            operation: operation.to_owned(),
316            source,
317        }),
318        Err(
319            crate::container::BoundedError::Stdin(error)
320            | crate::container::BoundedError::Wait(error),
321        ) => Err(ProcessCommandError::Io {
322            operation: operation.to_owned(),
323            source: error,
324        }),
325        Err(crate::container::BoundedError::WaitCleanup {
326            source, cleanup, ..
327        }) => Err(ProcessCommandError::WaitCleanup {
328            operation: operation.to_owned(),
329            source,
330            cleanup: cleanup.error.unwrap_or_else(|| {
331                if cleanup.verified {
332                    "verified"
333                } else {
334                    "unverified"
335                }
336                .to_owned()
337            }),
338        }),
339    }
340}
341
342impl ProcessObserver for SystemProcessRuntime {
343    fn status(&self, handle: &ProcessHandle) -> ProcessStatus {
344        match handle {
345            ProcessHandle::Local(handle) => verified_local_status(handle),
346            ProcessHandle::Ssh(handle) => verified_ssh_status(handle),
347        }
348    }
349
350    fn status_with_bound(&self, handle: &ProcessHandle, bound: &OperationBound) -> ProcessStatus {
351        match handle {
352            ProcessHandle::Local(handle) => verified_local_status_with_bound(handle, bound),
353            ProcessHandle::Ssh(handle) => verified_ssh_status_with_bound(handle, bound),
354        }
355    }
356
357    fn sync_logs(
358        &self,
359        handle: &ProcessHandle,
360        stdout: &Path,
361        stderr: &Path,
362        cleanup: bool,
363    ) -> Result<(), LogSyncError> {
364        match handle {
365            ProcessHandle::Local(_) => Ok(()),
366            ProcessHandle::Ssh(handle) => {
367                let bound = OperationBound::finite(REMOTE_LOG_SYNC_DEADLINE);
368                fetch_remote_file(&handle.target, &handle.stdout, stdout, &bound, cleanup)?;
369                fetch_remote_file(&handle.target, &handle.stderr, stderr, &bound, cleanup)
370            }
371        }
372    }
373}