//! PID reservations serialize detached launches; process identities protect stale-record cleanup. #[path = "pid_identity.rs"] 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; #[cfg(any(unix, windows))] 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; #[cfg(unix)] 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; #[derive(Debug)] #[cfg_attr(not(any(unix, windows)), allow(dead_code))] pub(crate) struct PidBackend { codex_bin: PathBuf, pid_file: PathBuf, lock_file: PathBuf, command_kind: PidCommandKind, } #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] struct PidRecord { pid: u32, // Keep the legacy timestamp for older CLI/updater versions reading this record. process_start_time: String, #[serde(default, skip_serializing_if = "Option::is_none")] #[serde(alias = "linuxProcessIdentity")] process_identity: Option, #[serde(default, skip_serializing_if = "Option::is_none")] executable_identity: Option, } #[derive(Debug, Clone, PartialEq, Eq)] 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); } } } #[derive(Debug, Clone, PartialEq, Eq)] enum PidFileState { Missing, Starting, Running(PidRecord), } #[derive(Debug, Clone)] #[cfg_attr(not(any(unix, windows)), allow(dead_code))] enum PidCommandKind { AppServer { remote_control_enabled: bool }, UpdateLoop { restore_release: Option }, } impl PidBackend { pub(crate) async fn running_executable_identity(&self) -> Result> { 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, ) -> 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 { 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, } } } } } #[cfg(any(unix, windows))] pub(crate) async fn start(&self) -> Result> { self.start_inner(/*replacement*/ None).await } #[cfg(not(any(unix, windows)))] pub(crate) async fn start(&self) -> Result> { 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; #[cfg(unix)] self.terminate_process(pid)?; #[cfg(windows)] 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 { #[cfg(unix)] 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 { #[cfg(unix)] self.force_terminate_process(pid)?; #[cfg(windows)] 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> { 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 { 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 { 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 { 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 { 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) } #[cfg(any(unix, windows))] async fn open_stderr_log(&self) -> Result { 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() ) }) } #[cfg(any(unix, windows))] 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 } } } #[cfg(any(unix, windows))] 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), #[cfg(unix)] PidCommandKind::UpdateLoop { .. } => terminate_process_group(pid), #[cfg(not(unix))] PidCommandKind::UpdateLoop { .. } => terminate_process(pid), } } #[cfg(not(windows))] 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 { process_matches_record(record).await } } pub(crate) async fn read_stderr_log_tail(pid_file: &Path) -> Result> { 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> { 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)) } #[cfg(unix)] 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) } #[cfg(unix)] 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}")) } #[cfg(unix)] 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}")) } #[cfg(unix)] 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}")) } #[cfg(unix)] 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}")) } #[cfg(not(any(unix, windows)))] fn terminate_process(_pid: u32) -> Result<()> { bail!("pid-managed app-server shutdown is unsupported on this platform") } #[cfg(not(any(unix, windows)))] fn force_terminate_process(_pid: u32) -> Result<()> { bail!("pid-managed app-server shutdown is unsupported on this platform") } #[cfg(not(any(unix, windows)))] fn force_terminate_process_group(_pid: u32) -> Result<()> { bail!("pid-managed updater shutdown is unsupported on this platform") } #[cfg(unix)] async fn process_matches_record(record: &PidRecord) -> Result { if !process_exists(record.pid) { return Ok(false); } let details = async { #[cfg(any(target_os = "linux", target_os = "macos"))] 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), } } #[cfg(not(any(unix, windows)))] async fn process_matches_record(_record: &PidRecord) -> Result { Ok(false) } #[derive(Debug, Clone, PartialEq, Eq)] #[cfg_attr(not(any(unix, windows)), allow(dead_code))] enum EmptyPidReservation { Active, Stale, Record(PidRecord), } #[cfg(unix)] fn try_lock_file(file: &fs::File) -> Result { 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") } #[cfg(not(any(unix, windows)))] fn try_lock_file(_file: &fs::File) -> Result { bail!("pid-managed app-server startup is unsupported on this platform") } #[cfg(any(unix, windows))] async fn reservation_lock_is_active(path: &Path) -> Result { 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)?) } #[cfg(not(any(unix, windows)))] async fn reservation_lock_is_active(_path: &Path) -> Result { Ok(false) } #[cfg(any(unix, windows))] async fn inspect_empty_pid_reservation( pid_path: &Path, lock_path: &Path, ) -> Result { 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)) } #[cfg(not(any(unix, windows)))] async fn inspect_empty_pid_reservation( _pid_path: &Path, _lock_path: &Path, ) -> Result { Ok(EmptyPidReservation::Stale) } #[cfg(unix)] async fn read_process_start_time(pid: u32) -> Result { Ok(read_process_details(pid).await?.1) } #[cfg(unix)] 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())) } #[cfg(all(test, any(unix, windows)))] #[path = "pid_tests.rs"] mod tests; #[cfg(windows)] use super::windows::try_lock_file; #[cfg(windows)] async fn read_process_start_time(pid: u32) -> Result { super::windows::Process::open(pid)? .context("daemon process exited during startup")? .start_time() } #[cfg(windows)] async fn process_matches_record(record: &PidRecord) -> Result { let Some(process) = super::windows::Process::open(record.pid)? else { return Ok(false); }; Ok(process.is_running()? && process.start_time()? == record.process_start_time) } #[cfg(windows)] fn force_terminate_process(pid: u32) -> Result<()> { if let Some(process) = super::windows::Process::open(pid)? { process.terminate()?; } Ok(()) } #[cfg(windows)] use force_terminate_process as terminate_process; #[cfg(windows)] #[path = "pid_windows.rs"] mod windows; #[cfg(any(unix, windows))] #[path = "pid_start.rs"] mod start;