Download codex-rs/app-server-daemon/src/backend/pid.rs from SaylorTwift/codex: direct link, hf CLI and curl.
- Browser
- Download file 26.8 kB
-
https://huggingface.co/SaylorTwift/codex/resolve/main/codex-rs/app-server-daemon/src/backend/pid.rs
- Command line
-
hf download hf://SaylorTwift/codex/codex-rs/app-server-daemon/src/backend/pid.rs
-
curl -L -o pid.rs https://huggingface.co/SaylorTwift/codex/resolve/main/codex-rs/app-server-daemon/src/backend/pid.rs
26.8 kB
| //! PID reservations serialize detached launches; process identities protect stale-record cleanup. | |
| mod identity; | |
| use std::io::SeekFrom; | |
| use std::path::Path; | |
| use std::path::PathBuf; | |
| use std::time::Duration; | |
| use crate::managed_install::ExecutableIdentity; | |
| use anyhow::Context; | |
| use anyhow::Result; | |
| use anyhow::bail; | |
| use codex_app_server_transport::REMOTE_CONTROL_DISABLED_ENV_VAR; | |
| use serde::Deserialize; | |
| use serde::Serialize; | |
| use tokio::fs; | |
| use tokio::io::AsyncReadExt; | |
| use tokio::io::AsyncSeekExt; | |
| use tokio::process::Command; | |
| use tokio::time::sleep; | |
| use crate::settings::DEFAULT_SHUTDOWN_GRACE_SECONDS; | |
| const STOP_POLL_INTERVAL: Duration = Duration::from_millis(50); | |
| const STOP_FORCE_TIMEOUT: Duration = Duration::from_secs(10); | |
| const START_TIMEOUT: Duration = Duration::from_secs(10); | |
| const STDERR_LOG_TAIL_BYTES: u64 = 4096; | |
| pub(crate) struct PidBackend { | |
| codex_bin: PathBuf, | |
| pid_file: PathBuf, | |
| lock_file: PathBuf, | |
| command_kind: PidCommandKind, | |
| } | |
| struct PidRecord { | |
| pid: u32, | |
| // Keep the legacy timestamp for older CLI/updater versions reading this record. | |
| process_start_time: String, | |
| process_identity: Option<identity::ProcessIdentity>, | |
| executable_identity: Option<ExecutableIdentity>, | |
| } | |
| pub(crate) struct PidLogTail { | |
| pub(crate) path: PathBuf, | |
| pub(crate) contents: String, | |
| } | |
| impl PidLogTail { | |
| pub(crate) fn append_to_context(&self, context: &mut String) { | |
| context.push_str(&format!( | |
| "\n\nManaged app-server stderr ({}):", | |
| self.path.display() | |
| )); | |
| for line in self.contents.lines() { | |
| context.push_str("\n "); | |
| context.push_str(line); | |
| } | |
| } | |
| } | |
| enum PidFileState { | |
| Missing, | |
| Starting, | |
| Running(PidRecord), | |
| } | |
| enum PidCommandKind { | |
| AppServer { remote_control_enabled: bool }, | |
| UpdateLoop { restore_release: Option<String> }, | |
| } | |
| impl PidBackend { | |
| pub(crate) async fn running_executable_identity(&self) -> Result<Option<ExecutableIdentity>> { | |
| match self.read_pid_file_state().await? { | |
| PidFileState::Running(record) if self.record_is_active(&record).await? => { | |
| Ok(record.executable_identity) | |
| } | |
| _ => Ok(None), | |
| } | |
| } | |
| pub(crate) fn new(codex_bin: PathBuf, pid_file: PathBuf, remote_control_enabled: bool) -> Self { | |
| let lock_file = pid_file.with_extension("pid.lock"); | |
| Self { | |
| codex_bin, | |
| pid_file, | |
| lock_file, | |
| command_kind: PidCommandKind::AppServer { | |
| remote_control_enabled, | |
| }, | |
| } | |
| } | |
| pub(crate) fn new_update_loop( | |
| codex_bin: PathBuf, | |
| pid_file: PathBuf, | |
| restore_release: Option<String>, | |
| ) -> Self { | |
| let lock_file = pid_file.with_extension("pid.lock"); | |
| Self { | |
| codex_bin, | |
| pid_file, | |
| lock_file, | |
| command_kind: PidCommandKind::UpdateLoop { restore_release }, | |
| } | |
| } | |
| pub(crate) async fn is_starting_or_running(&self) -> Result<bool> { | |
| loop { | |
| match self.read_pid_file_state().await? { | |
| PidFileState::Missing => return Ok(false), | |
| PidFileState::Starting => return Ok(true), | |
| PidFileState::Running(record) => { | |
| if self.record_is_active(&record).await? { | |
| return Ok(true); | |
| } | |
| match self.refresh_after_stale_record(&record).await? { | |
| PidFileState::Missing => return Ok(false), | |
| PidFileState::Starting | PidFileState::Running(_) => continue, | |
| } | |
| } | |
| } | |
| } | |
| } | |
| pub(crate) async fn start(&self) -> Result<Option<u32>> { | |
| self.start_inner(/*replacement*/ None).await | |
| } | |
| pub(crate) async fn start(&self) -> Result<Option<u32>> { | |
| bail!("pid-managed app-server startup is unsupported on this platform") | |
| } | |
| pub(crate) async fn stop(&self) -> Result<()> { | |
| self.stop_with_grace(DEFAULT_SHUTDOWN_GRACE_SECONDS).await | |
| } | |
| pub(crate) async fn stop_with_grace(&self, grace_seconds: u32) -> Result<()> { | |
| loop { | |
| let Some(record) = self.wait_for_pid_start().await? else { | |
| return Ok(()); | |
| }; | |
| if !self.record_is_active(&record).await? { | |
| match self.refresh_after_stale_record(&record).await? { | |
| PidFileState::Missing => return Ok(()), | |
| PidFileState::Starting | PidFileState::Running(_) => continue, | |
| } | |
| } | |
| let pid = record.pid; | |
| let started_at = tokio::time::Instant::now(); | |
| let force_after = Duration::from_secs(grace_seconds.into()); | |
| let deadline = started_at + force_after + STOP_FORCE_TIMEOUT; | |
| self.terminate_process(pid)?; | |
| let process = { | |
| let Some(process) = super::windows::Process::open(pid)? else { | |
| continue; | |
| }; | |
| if process.start_time()? != record.process_start_time { | |
| continue; | |
| } | |
| match self.command_kind { | |
| PidCommandKind::AppServer { .. } => { | |
| let codex_home = self | |
| .pid_file | |
| .parent() | |
| .and_then(Path::parent) | |
| .context("daemon pid path has no Codex home")?; | |
| let socket_path = | |
| codex_app_server_transport::app_server_control_socket_path(codex_home)?; | |
| if let Err(err) = | |
| crate::client::request_shutdown(socket_path.as_path(), pid).await | |
| { | |
| tracing::warn!(%pid, %err, "managed app-server shutdown request failed; waiting for force deadline"); | |
| } | |
| } | |
| PidCommandKind::UpdateLoop { .. } => { | |
| fs::write(self.pid_file.with_extension("shutdown"), pid.to_string()) | |
| .await | |
| .context("failed to request updater shutdown")?; | |
| } | |
| } | |
| process | |
| }; | |
| let mut forced = false; | |
| loop { | |
| if let Ok(raw_pid) = libc::pid_t::try_from(pid) | |
| && raw_pid > 0 | |
| { | |
| // A previous updater may have started this child; reap it if it has exited. | |
| unsafe { libc::waitpid(raw_pid, std::ptr::null_mut(), libc::WNOHANG) }; | |
| } | |
| if !self.record_is_active(&record).await? { | |
| match self.refresh_after_stale_record(&record).await? { | |
| PidFileState::Missing => return Ok(()), | |
| PidFileState::Starting | PidFileState::Running(_) => break, | |
| } | |
| } | |
| if tokio::time::Instant::now() >= deadline { | |
| break; | |
| } | |
| if !forced && started_at.elapsed() >= force_after { | |
| self.force_terminate_process(pid)?; | |
| process.terminate()?; | |
| forced = true; | |
| } | |
| sleep(STOP_POLL_INTERVAL).await; | |
| } | |
| if self.record_is_active(&record).await? { | |
| bail!("timed out waiting for pid-managed app server {pid} to stop"); | |
| } | |
| } | |
| } | |
| async fn wait_for_pid_start(&self) -> Result<Option<PidRecord>> { | |
| let deadline = tokio::time::Instant::now() + START_TIMEOUT; | |
| loop { | |
| match self.read_pid_file_state().await? { | |
| PidFileState::Missing => return Ok(None), | |
| PidFileState::Running(record) => return Ok(Some(record)), | |
| PidFileState::Starting if tokio::time::Instant::now() < deadline => { | |
| sleep(STOP_POLL_INTERVAL).await; | |
| } | |
| PidFileState::Starting => { | |
| bail!( | |
| "timed out waiting for pid reservation in {} to finish initializing", | |
| self.pid_file.display() | |
| ); | |
| } | |
| } | |
| } | |
| } | |
| async fn read_pid_file_state(&self) -> Result<PidFileState> { | |
| let contents = match fs::read_to_string(&self.pid_file).await { | |
| Ok(contents) => contents, | |
| Err(err) if err.kind() == std::io::ErrorKind::NotFound => { | |
| return if reservation_lock_is_active(&self.lock_file).await? { | |
| Ok(PidFileState::Starting) | |
| } else { | |
| Ok(PidFileState::Missing) | |
| }; | |
| } | |
| Err(err) => { | |
| return Err(err).with_context(|| { | |
| format!("failed to read pid file {}", self.pid_file.display()) | |
| }); | |
| } | |
| }; | |
| if contents.trim().is_empty() { | |
| match inspect_empty_pid_reservation(&self.pid_file, &self.lock_file).await? { | |
| EmptyPidReservation::Active => { | |
| return Ok(PidFileState::Starting); | |
| } | |
| EmptyPidReservation::Stale => { | |
| return Ok(PidFileState::Missing); | |
| } | |
| EmptyPidReservation::Record(record) => return Ok(PidFileState::Running(record)), | |
| } | |
| } | |
| let record = serde_json::from_str(&contents) | |
| .with_context(|| format!("invalid pid file contents in {}", self.pid_file.display()))?; | |
| Ok(PidFileState::Running(record)) | |
| } | |
| async fn read_pid_file_state_with_lock_held(&self) -> Result<PidFileState> { | |
| let contents = match fs::read_to_string(&self.pid_file).await { | |
| Ok(contents) => contents, | |
| Err(err) if err.kind() == std::io::ErrorKind::NotFound => { | |
| return Ok(PidFileState::Missing); | |
| } | |
| Err(err) => { | |
| return Err(err).with_context(|| { | |
| format!("failed to read pid file {}", self.pid_file.display()) | |
| }); | |
| } | |
| }; | |
| if contents.trim().is_empty() { | |
| let _ = fs::remove_file(&self.pid_file).await; | |
| return Ok(PidFileState::Missing); | |
| } | |
| let record = serde_json::from_str(&contents) | |
| .with_context(|| format!("invalid pid file contents in {}", self.pid_file.display()))?; | |
| Ok(PidFileState::Running(record)) | |
| } | |
| async fn refresh_after_stale_record(&self, expected: &PidRecord) -> Result<PidFileState> { | |
| let reservation_lock = self.acquire_reservation_lock().await?; | |
| let state = match self.read_pid_file_state_with_lock_held().await? { | |
| PidFileState::Running(record) if record == *expected => { | |
| let _ = fs::remove_file(&self.pid_file).await; | |
| PidFileState::Missing | |
| } | |
| state => state, | |
| }; | |
| drop(reservation_lock); | |
| Ok(state) | |
| } | |
| async fn acquire_reservation_lock(&self) -> Result<fs::File> { | |
| let reservation_lock = fs::OpenOptions::new() | |
| .create(true) | |
| .truncate(false) | |
| .write(true) | |
| .open(&self.lock_file) | |
| .await | |
| .with_context(|| { | |
| format!("failed to open pid lock file {}", self.lock_file.display()) | |
| })?; | |
| let lock_deadline = tokio::time::Instant::now() + START_TIMEOUT; | |
| while !try_lock_file(&reservation_lock)? { | |
| if tokio::time::Instant::now() >= lock_deadline { | |
| bail!( | |
| "timed out waiting for pid lock {}", | |
| self.lock_file.display() | |
| ); | |
| } | |
| sleep(STOP_POLL_INTERVAL).await; | |
| } | |
| Ok(reservation_lock) | |
| } | |
| async fn open_stderr_log(&self) -> Result<fs::File> { | |
| let stderr_log_file = stderr_log_file_for_pid_file(&self.pid_file); | |
| fs::OpenOptions::new() | |
| .create(true) | |
| .truncate(true) | |
| .write(true) | |
| .open(&stderr_log_file) | |
| .await | |
| .with_context(|| { | |
| format!( | |
| "failed to open stderr log for pid-managed app server {}", | |
| stderr_log_file.display() | |
| ) | |
| }) | |
| } | |
| fn command_args(&self) -> Vec<&str> { | |
| match &self.command_kind { | |
| PidCommandKind::AppServer { | |
| remote_control_enabled: true, | |
| } => vec!["app-server", "--remote-control", "--listen", "unix://"], | |
| PidCommandKind::AppServer { | |
| remote_control_enabled: false, | |
| } => vec!["app-server", "--listen", "unix://"], | |
| PidCommandKind::UpdateLoop { restore_release } => { | |
| let mut args = vec!["app-server", "daemon", "pid-update-loop"]; | |
| if let Some(release) = restore_release { | |
| args.extend(["--restore-release", release.as_str()]); | |
| } | |
| args | |
| } | |
| } | |
| } | |
| fn command_env(&self) -> Option<(&'static str, &'static str)> { | |
| match self.command_kind { | |
| PidCommandKind::AppServer { | |
| remote_control_enabled: false, | |
| } => Some((REMOTE_CONTROL_DISABLED_ENV_VAR, "1")), | |
| PidCommandKind::AppServer { | |
| remote_control_enabled: true, | |
| } | |
| | PidCommandKind::UpdateLoop { .. } => None, | |
| } | |
| } | |
| fn terminate_process(&self, pid: u32) -> Result<()> { | |
| match self.command_kind { | |
| PidCommandKind::AppServer { .. } => terminate_process(pid), | |
| PidCommandKind::UpdateLoop { .. } => terminate_process_group(pid), | |
| PidCommandKind::UpdateLoop { .. } => terminate_process(pid), | |
| } | |
| } | |
| fn force_terminate_process(&self, pid: u32) -> Result<()> { | |
| match self.command_kind { | |
| PidCommandKind::AppServer { .. } => force_terminate_process(pid), | |
| PidCommandKind::UpdateLoop { .. } => force_terminate_process_group(pid), | |
| } | |
| } | |
| async fn record_is_active(&self, record: &PidRecord) -> Result<bool> { | |
| process_matches_record(record).await | |
| } | |
| } | |
| pub(crate) async fn read_stderr_log_tail(pid_file: &Path) -> Result<Option<PidLogTail>> { | |
| let path = stderr_log_file_for_pid_file(pid_file); | |
| let Some(contents) = read_log_tail(&path, STDERR_LOG_TAIL_BYTES).await? else { | |
| return Ok(None); | |
| }; | |
| Ok(Some(PidLogTail { path, contents })) | |
| } | |
| fn stderr_log_file_for_pid_file(pid_file: &Path) -> PathBuf { | |
| pid_file.with_extension("stderr.log") | |
| } | |
| async fn read_log_tail(path: &Path, byte_limit: u64) -> Result<Option<String>> { | |
| let mut file = match fs::File::open(path).await { | |
| Ok(file) => file, | |
| Err(err) if err.kind() == std::io::ErrorKind::NotFound => return Ok(None), | |
| Err(err) => { | |
| return Err(err) | |
| .with_context(|| format!("failed to open stderr log {}", path.display())); | |
| } | |
| }; | |
| let len = file | |
| .metadata() | |
| .await | |
| .with_context(|| format!("failed to inspect stderr log {}", path.display()))? | |
| .len(); | |
| if len == 0 { | |
| return Ok(None); | |
| } | |
| let start = len.saturating_sub(byte_limit); | |
| file.seek(SeekFrom::Start(start)) | |
| .await | |
| .with_context(|| format!("failed to seek stderr log {}", path.display()))?; | |
| let mut bytes = Vec::new(); | |
| file.read_to_end(&mut bytes) | |
| .await | |
| .with_context(|| format!("failed to read stderr log {}", path.display()))?; | |
| if start > 0 | |
| && let Some(newline_index) = bytes.iter().position(|byte| *byte == b'\n') | |
| { | |
| bytes.drain(..=newline_index); | |
| } | |
| let contents = String::from_utf8_lossy(&bytes).trim_end().to_string(); | |
| if contents.is_empty() { | |
| return Ok(None); | |
| } | |
| Ok(Some(contents)) | |
| } | |
| fn process_exists(pid: u32) -> bool { | |
| let Ok(pid) = libc::pid_t::try_from(pid) else { | |
| return false; | |
| }; | |
| let result = unsafe { libc::kill(pid, 0) }; | |
| result == 0 || std::io::Error::last_os_error().raw_os_error() == Some(libc::EPERM) | |
| } | |
| fn terminate_process(pid: u32) -> Result<()> { | |
| let raw_pid = libc::pid_t::try_from(pid) | |
| .with_context(|| format!("pid-managed app server pid {pid} is out of range"))?; | |
| let result = unsafe { libc::kill(raw_pid, libc::SIGTERM) }; | |
| if result == 0 { | |
| return Ok(()); | |
| } | |
| let err = std::io::Error::last_os_error(); | |
| if err.raw_os_error() == Some(libc::ESRCH) { | |
| return Ok(()); | |
| } | |
| Err(err).with_context(|| format!("failed to terminate pid-managed app server {pid}")) | |
| } | |
| fn force_terminate_process(pid: u32) -> Result<()> { | |
| let raw_pid = libc::pid_t::try_from(pid) | |
| .with_context(|| format!("pid-managed app server pid {pid} is out of range"))?; | |
| let result = unsafe { libc::kill(raw_pid, libc::SIGKILL) }; | |
| if result == 0 { | |
| return Ok(()); | |
| } | |
| let err = std::io::Error::last_os_error(); | |
| if err.raw_os_error() == Some(libc::ESRCH) { | |
| return Ok(()); | |
| } | |
| Err(err).with_context(|| format!("failed to force terminate pid-managed app server {pid}")) | |
| } | |
| fn terminate_process_group(pid: u32) -> Result<()> { | |
| let raw_pid = libc::pid_t::try_from(pid) | |
| .with_context(|| format!("pid-managed updater pid {pid} is out of range"))?; | |
| let result = unsafe { libc::kill(-raw_pid, libc::SIGTERM) }; | |
| if result == 0 { | |
| return Ok(()); | |
| } | |
| let err = std::io::Error::last_os_error(); | |
| if err.raw_os_error() == Some(libc::ESRCH) { | |
| return Ok(()); | |
| } | |
| Err(err).with_context(|| format!("failed to terminate pid-managed updater group {pid}")) | |
| } | |
| fn force_terminate_process_group(pid: u32) -> Result<()> { | |
| let raw_pid = libc::pid_t::try_from(pid) | |
| .with_context(|| format!("pid-managed updater pid {pid} is out of range"))?; | |
| let result = unsafe { libc::kill(-raw_pid, libc::SIGKILL) }; | |
| if result == 0 { | |
| return Ok(()); | |
| } | |
| let err = std::io::Error::last_os_error(); | |
| if err.raw_os_error() == Some(libc::ESRCH) { | |
| return Ok(()); | |
| } | |
| Err(err).with_context(|| format!("failed to force terminate pid-managed updater group {pid}")) | |
| } | |
| fn terminate_process(_pid: u32) -> Result<()> { | |
| bail!("pid-managed app-server shutdown is unsupported on this platform") | |
| } | |
| fn force_terminate_process(_pid: u32) -> Result<()> { | |
| bail!("pid-managed app-server shutdown is unsupported on this platform") | |
| } | |
| fn force_terminate_process_group(_pid: u32) -> Result<()> { | |
| bail!("pid-managed updater shutdown is unsupported on this platform") | |
| } | |
| async fn process_matches_record(record: &PidRecord) -> Result<bool> { | |
| if !process_exists(record.pid) { | |
| return Ok(false); | |
| } | |
| let details = async { | |
| if let Some(expected) = &record.process_identity { | |
| return expected.matches_process(record.pid).await; | |
| } | |
| let (state, start_time) = read_process_details(record.pid).await?; | |
| let matches = start_time == record.process_start_time; | |
| let is_zombie = state.starts_with('Z'); | |
| if !matches && !is_zombie { | |
| bail!( | |
| "cannot verify pid-managed process {}: legacy start time changed; PID record retained. \ | |
| Retry with the locale and timezone used to start the daemon. If the system clock \ | |
| changed, stop the original process before restarting the daemon", | |
| record.pid | |
| ); | |
| } | |
| Ok((is_zombie, matches)) | |
| } | |
| .await; | |
| match details { | |
| Ok((is_zombie, matches)) => { | |
| // An unreaped zombie still passes kill(pid, 0) and retains its start | |
| // time, but it can no longer run the app-server or updater. | |
| if is_zombie { | |
| if matches && let Ok(raw_pid) = libc::pid_t::try_from(record.pid) { | |
| // Re-exec can lose the Child handle without changing parenthood. | |
| unsafe { libc::waitpid(raw_pid, std::ptr::null_mut(), libc::WNOHANG) }; | |
| } | |
| return Ok(false); | |
| } | |
| Ok(matches) | |
| } | |
| Err(_err) if !process_exists(record.pid) => Ok(false), | |
| Err(err) => Err(err), | |
| } | |
| } | |
| async fn process_matches_record(_record: &PidRecord) -> Result<bool> { | |
| Ok(false) | |
| } | |
| enum EmptyPidReservation { | |
| Active, | |
| Stale, | |
| Record(PidRecord), | |
| } | |
| fn try_lock_file(file: &fs::File) -> Result<bool> { | |
| use std::os::fd::AsRawFd; | |
| let result = unsafe { libc::flock(file.as_raw_fd(), libc::LOCK_EX | libc::LOCK_NB) }; | |
| if result == 0 { | |
| return Ok(true); | |
| } | |
| let err = std::io::Error::last_os_error(); | |
| if err.raw_os_error() == Some(libc::EWOULDBLOCK) { | |
| return Ok(false); | |
| } | |
| Err(err).context("failed to lock pid reservation") | |
| } | |
| fn try_lock_file(_file: &fs::File) -> Result<bool> { | |
| bail!("pid-managed app-server startup is unsupported on this platform") | |
| } | |
| async fn reservation_lock_is_active(path: &Path) -> Result<bool> { | |
| let file = match fs::OpenOptions::new() | |
| .write(true) | |
| .create(true) | |
| .truncate(false) | |
| .open(path) | |
| .await | |
| { | |
| Ok(file) => file, | |
| Err(err) if err.kind() == std::io::ErrorKind::NotFound => { | |
| return Ok(false); | |
| } | |
| Err(err) => { | |
| return Err(err) | |
| .with_context(|| format!("failed to inspect pid lock file {}", path.display())); | |
| } | |
| }; | |
| Ok(!try_lock_file(&file)?) | |
| } | |
| async fn reservation_lock_is_active(_path: &Path) -> Result<bool> { | |
| Ok(false) | |
| } | |
| async fn inspect_empty_pid_reservation( | |
| pid_path: &Path, | |
| lock_path: &Path, | |
| ) -> Result<EmptyPidReservation> { | |
| let file = match fs::OpenOptions::new() | |
| .write(true) | |
| .create(true) | |
| .truncate(false) | |
| .open(lock_path) | |
| .await | |
| { | |
| Ok(file) => file, | |
| Err(err) => { | |
| return Err(err).with_context(|| { | |
| format!("failed to inspect pid lock file {}", lock_path.display()) | |
| }); | |
| } | |
| }; | |
| if !try_lock_file(&file)? { | |
| return Ok(EmptyPidReservation::Active); | |
| } | |
| let contents = match fs::read_to_string(pid_path).await { | |
| Ok(contents) => contents, | |
| Err(err) if err.kind() == std::io::ErrorKind::NotFound => { | |
| return Ok(EmptyPidReservation::Stale); | |
| } | |
| Err(err) => { | |
| return Err(err) | |
| .with_context(|| format!("failed to reread pid file {}", pid_path.display())); | |
| } | |
| }; | |
| if contents.trim().is_empty() { | |
| let _ = fs::remove_file(pid_path).await; | |
| return Ok(EmptyPidReservation::Stale); | |
| } | |
| let record = serde_json::from_str(&contents) | |
| .with_context(|| format!("invalid pid file contents in {}", pid_path.display()))?; | |
| Ok(EmptyPidReservation::Record(record)) | |
| } | |
| async fn inspect_empty_pid_reservation( | |
| _pid_path: &Path, | |
| _lock_path: &Path, | |
| ) -> Result<EmptyPidReservation> { | |
| Ok(EmptyPidReservation::Stale) | |
| } | |
| async fn read_process_start_time(pid: u32) -> Result<String> { | |
| Ok(read_process_details(pid).await?.1) | |
| } | |
| async fn read_process_details(pid: u32) -> Result<(String, String)> { | |
| let output = Command::new("ps") | |
| .args(["-p", &pid.to_string(), "-o", "stat=", "-o", "lstart="]) | |
| .output() | |
| .await | |
| .context("failed to invoke ps for pid-managed app server")?; | |
| if !output.status.success() { | |
| bail!("failed to read start time for pid-managed app server {pid}"); | |
| } | |
| let details = String::from_utf8(output.stdout) | |
| .context("pid-managed app server process details were not utf-8")?; | |
| let Some((state, start_time)) = details.trim().split_once(char::is_whitespace) else { | |
| bail!("pid-managed app server {pid} has no recorded start time"); | |
| }; | |
| let start_time = start_time.trim(); | |
| if start_time.is_empty() { | |
| bail!("pid-managed app server {pid} has no recorded start time"); | |
| } | |
| Ok((state.to_string(), start_time.to_string())) | |
| } | |
| mod tests; | |
| use super::windows::try_lock_file; | |
| async fn read_process_start_time(pid: u32) -> Result<String> { | |
| super::windows::Process::open(pid)? | |
| .context("daemon process exited during startup")? | |
| .start_time() | |
| } | |
| async fn process_matches_record(record: &PidRecord) -> Result<bool> { | |
| let Some(process) = super::windows::Process::open(record.pid)? else { | |
| return Ok(false); | |
| }; | |
| Ok(process.is_running()? && process.start_time()? == record.process_start_time) | |
| } | |
| fn force_terminate_process(pid: u32) -> Result<()> { | |
| if let Some(process) = super::windows::Process::open(pid)? { | |
| process.terminate()?; | |
| } | |
| Ok(()) | |
| } | |
| use force_terminate_process as terminate_process; | |
| mod windows; | |
| mod start; | |