use std::sync::mpsc::{self, RecvTimeoutError, Sender};
use std::time::Duration;
pub const WATCHDOG_GRACE: Duration = Duration::from_secs(20);
pub struct ConnectWatchdog {
disarm: Option<Sender<()>>,
}
impl ConnectWatchdog {
pub fn arm(hard_cap: Duration) -> Self {
let (tx, rx) = mpsc::channel::<()>();
let _ = std::thread::Builder::new()
.name("rp-connect-watchdog".to_owned())
.spawn(move || match rx.recv_timeout(hard_cap) {
Ok(()) | Err(RecvTimeoutError::Disconnected) => {}
Err(RecvTimeoutError::Timeout) => {
eprintln!(
"running-process: connect watchdog fired after {hard_cap:?} \
without a reachable daemon or a clean exit; aborting to \
avoid a hang (running-process#894)"
);
std::process::abort();
}
});
Self { disarm: Some(tx) }
}
}
impl Drop for ConnectWatchdog {
fn drop(&mut self) {
if let Some(tx) = self.disarm.take() {
let _ = tx.send(());
}
}
}
pub fn capture_connect_dump(
program: &str,
deadline: Duration,
error: &str,
) -> Option<std::path::PathBuf> {
#[cfg(feature = "probe")]
{
use running_process_probe::snapshot::{capture_and_resolve, SnapshotConfig};
let snapshot = capture_and_resolve(&SnapshotConfig::default()).ok()?;
if snapshot.threads.is_empty() {
return None;
}
let report = render_connect_dump(program, deadline, error, &snapshot);
eprint!("{report}");
let path = std::env::temp_dir().join(format!(
"rp-connect-dump-{}-{}.txt",
sanitize(program),
std::process::id()
));
std::fs::write(&path, &report).ok().map(|()| path)
}
#[cfg(not(feature = "probe"))]
{
let _ = (program, deadline, error);
None
}
}
#[cfg(feature = "probe")]
fn sanitize(program: &str) -> String {
program
.chars()
.map(|c| if c.is_ascii_alphanumeric() { c } else { '-' })
.collect()
}
#[cfg(feature = "probe")]
fn render_connect_dump(
program: &str,
deadline: Duration,
error: &str,
snapshot: &running_process_probe::snapshot::Snapshot,
) -> String {
use std::fmt::Write as _;
let mut out = String::new();
let _ = writeln!(
out,
"running-process: v2 broker for '{program}' unreachable within {deadline:?}: {error}"
);
let _ = writeln!(
out,
"all-thread stack dump ({} sibling thread(s), frames_resolved={}):",
snapshot.threads.len(),
snapshot.frames_resolved
);
for (idx, thread) in snapshot.threads.iter().enumerate() {
let _ = writeln!(
out,
" thread #{idx} os_tid={} ip={:#018x}{}",
thread.os_tid,
thread.instruction_pointer,
if thread.truncated {
" (stack truncated)"
} else {
""
}
);
if thread.frames.is_empty() {
let _ = writeln!(out, " <no resolvable frames>");
}
for (fidx, addr) in thread.frames.iter().enumerate() {
let _ = writeln!(out, " #{fidx:<3} {addr:#018x}");
}
}
out
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn watchdog_disarms_on_drop_without_aborting() {
let guard = ConnectWatchdog::arm(Duration::from_secs(3600));
drop(guard);
std::thread::sleep(Duration::from_millis(50));
}
#[test]
fn dump_is_none_without_probe_feature() {
#[cfg(not(feature = "probe"))]
assert!(capture_connect_dump("zccache", Duration::from_secs(3), "not found").is_none());
#[cfg(feature = "probe")]
{
let _ = capture_connect_dump("zccache", Duration::from_secs(3), "not found");
}
}
}