1use std::path::{Path, PathBuf};
32use std::time::Duration;
33
34pub fn sweep(root: &Path, grace: Duration) {
47 #[cfg(unix)]
48 {
49 let deadline = std::time::Instant::now() + grace;
50 let mut termed = false;
51 loop {
52 let left = stragglers(root);
53 if left.is_empty() {
54 return;
55 }
56 if !termed {
57 eprintln!(
58 "bynk dev: stopping {} leftover process(es) under {}…",
59 left.len(),
60 root.display()
61 );
62 if !signal("TERM", &left) {
65 report_survivors("could not run `kill`", &left);
66 return;
67 }
68 termed = true;
69 } else if std::time::Instant::now() >= deadline {
70 signal("KILL", &left);
71 let soon = std::time::Instant::now() + Duration::from_secs(1);
74 let mut left = stragglers(root);
75 while !left.is_empty() && std::time::Instant::now() < soon {
76 std::thread::sleep(Duration::from_millis(20));
77 left = stragglers(root);
78 }
79 if !left.is_empty() {
80 report_survivors("they survived SIGKILL", &left);
81 }
82 return;
83 }
84 std::thread::sleep(Duration::from_millis(100));
85 }
86 }
87 #[cfg(not(unix))]
89 let _ = (root, grace);
90}
91
92#[cfg(unix)]
94fn report_survivors(why: &str, pids: &[u32]) {
95 let pids = pids
96 .iter()
97 .map(u32::to_string)
98 .collect::<Vec<_>>()
99 .join(" ");
100 eprintln!(
101 "bynk dev: could not stop leftover processes ({why}): {pids}. They may hold \
102 ports the next `bynk dev` needs; stop them by hand."
103 );
104}
105
106#[cfg(unix)]
108pub fn stragglers(root: &Path) -> Vec<u32> {
109 let Some(table) = process_table() else {
110 return Vec::new();
111 };
112 let root = root.canonicalize().unwrap_or_else(|_| root.to_path_buf());
115 let candidates = in_own_group(&table, std::process::id());
116 let cwds = cwds(&candidates);
117 candidates
118 .into_iter()
119 .filter(|pid| {
120 cwds.iter()
121 .any(|(p, cwd)| p == pid && cwd.starts_with(&root))
122 })
123 .collect()
124}
125
126#[derive(Debug, Clone, Copy, PartialEq, Eq)]
128#[cfg_attr(not(unix), allow(dead_code))]
129struct Row {
130 pid: u32,
131 ppid: u32,
132 pgid: u32,
133}
134
135#[cfg(unix)]
138fn process_table() -> Option<Vec<Row>> {
139 let out = std::process::Command::new("ps")
140 .args(["-A", "-o", "pid=,ppid=,pgid="])
141 .stderr(std::process::Stdio::null())
142 .output()
143 .ok()?;
144 out.status
145 .success()
146 .then(|| parse_ps(&String::from_utf8_lossy(&out.stdout)))
147}
148
149#[cfg_attr(not(unix), allow(dead_code))]
150fn parse_ps(text: &str) -> Vec<Row> {
151 text.lines()
152 .filter_map(|line| {
153 let mut fields = line.split_whitespace().map(str::parse::<u32>);
154 let row = Row {
155 pid: fields.next()?.ok()?,
156 ppid: fields.next()?.ok()?,
157 pgid: fields.next()?.ok()?,
158 };
159 Some(row)
160 })
161 .collect()
162}
163
164#[cfg_attr(not(unix), allow(dead_code))]
169fn in_own_group(table: &[Row], me: u32) -> Vec<u32> {
170 let Some(own) = table.iter().find(|r| r.pid == me) else {
171 return Vec::new();
172 };
173 let mut ancestors = vec![me];
174 let mut at = own.ppid;
175 for _ in 0..table.len() {
177 if at == 0 || ancestors.contains(&at) {
178 break;
179 }
180 ancestors.push(at);
181 match table.iter().find(|r| r.pid == at) {
182 Some(r) => at = r.ppid,
183 None => break,
184 }
185 }
186 table
187 .iter()
188 .filter(|r| r.pgid == own.pgid && !ancestors.contains(&r.pid))
189 .map(|r| r.pid)
190 .collect()
191}
192
193#[cfg(target_os = "linux")]
195fn cwds(pids: &[u32]) -> Vec<(u32, PathBuf)> {
196 pids.iter()
197 .filter_map(|&pid| Some((pid, std::fs::read_link(format!("/proc/{pid}/cwd")).ok()?)))
198 .collect()
199}
200
201#[cfg(all(unix, not(target_os = "linux")))]
204fn cwds(pids: &[u32]) -> Vec<(u32, PathBuf)> {
205 if pids.is_empty() {
206 return Vec::new();
207 }
208 let list = pids
209 .iter()
210 .map(u32::to_string)
211 .collect::<Vec<_>>()
212 .join(",");
213 match std::process::Command::new("lsof")
217 .args(["-a", "-p", &list, "-d", "cwd", "-F", "pn"])
218 .stderr(std::process::Stdio::null())
219 .output()
220 {
221 Ok(out) => parse_lsof(&String::from_utf8_lossy(&out.stdout)),
222 Err(_) => Vec::new(),
223 }
224}
225
226#[cfg_attr(not(all(unix, not(target_os = "linux"))), allow(dead_code))]
230fn parse_lsof(text: &str) -> Vec<(u32, PathBuf)> {
231 let mut pid = None;
232 let mut found = Vec::new();
233 for line in text.lines() {
234 if let Some(p) = line.strip_prefix('p') {
235 pid = p.parse::<u32>().ok();
236 } else if let (Some(path), Some(pid)) = (line.strip_prefix('n'), pid) {
237 found.push((pid, PathBuf::from(path)));
238 }
239 }
240 found
241}
242
243#[cfg(unix)]
249fn signal(signal: &str, pids: &[u32]) -> bool {
250 std::process::Command::new("kill")
251 .arg(format!("-{signal}"))
252 .args(pids.iter().map(u32::to_string))
253 .stderr(std::process::Stdio::null())
254 .status()
255 .is_ok()
256}
257
258#[cfg(test)]
259mod tests {
260 use super::*;
261
262 fn row(pid: u32, ppid: u32, pgid: u32) -> Row {
263 Row { pid, ppid, pgid }
264 }
265
266 #[test]
267 fn ps_rows_parse_and_junk_lines_are_skipped() {
268 let text = " 1 0 1\n 200 1 200\nnot a row\n 201 200 200\n";
269 assert_eq!(
270 parse_ps(text),
271 vec![row(1, 0, 1), row(200, 1, 200), row(201, 200, 200)]
272 );
273 }
274
275 #[test]
279 fn the_group_leaves_out_self_and_ancestors_but_keeps_orphans() {
280 let table = [
281 row(1, 0, 1),
282 row(100, 1, 100), row(200, 100, 200), row(300, 200, 200), row(301, 200, 200), row(310, 300, 200), row(400, 1, 200), row(500, 1, 500), ];
290 let mut got = in_own_group(&table, 300);
291 got.sort_unstable();
292 assert_eq!(got, vec![301, 310, 400]);
295 }
296
297 #[test]
298 fn an_unlisted_self_selects_nothing() {
299 assert!(in_own_group(&[row(1, 0, 1)], 999).is_empty());
300 }
301
302 #[test]
303 fn lsof_output_parses_pid_and_cwd() {
304 let text = "p310\nfcwd\nn/private/tmp/p/.bynk/dev/workers/a\np400\nfcwd\nn/Users/me\n";
305 assert_eq!(
306 parse_lsof(text),
307 vec![
308 (310, PathBuf::from("/private/tmp/p/.bynk/dev/workers/a")),
309 (400, PathBuf::from("/Users/me")),
310 ]
311 );
312 }
313
314 #[cfg(unix)]
319 fn orphans_are_swept(tag: &str, script: &str, grace: Duration) {
320 use std::process::{Command, Stdio};
321 let root = std::env::temp_dir().join(format!("bynk-sweep-{tag}-{}", std::process::id()));
322 let _ = std::fs::remove_dir_all(&root);
323 let worker = root.join("workers/ctx");
324 std::fs::create_dir_all(&worker).unwrap();
325 let mut sh = Command::new("sh")
326 .args(["-c", script])
327 .current_dir(&worker)
328 .stdin(Stdio::null())
329 .stdout(Stdio::null())
330 .stderr(Stdio::null())
331 .spawn()
332 .expect("sh spawns");
333 let started = std::time::Instant::now();
334 while stragglers(&root).len() < 3 {
335 assert!(
336 started.elapsed() < Duration::from_secs(10),
337 "the script never started its background processes"
338 );
339 std::thread::sleep(Duration::from_millis(20));
340 }
341 let mut bystander = Command::new("sleep")
344 .arg("300")
345 .current_dir(std::env::temp_dir())
346 .stdin(Stdio::null())
347 .stdout(Stdio::null())
348 .stderr(Stdio::null())
349 .spawn()
350 .expect("sleep spawns");
351 let _ = sh.kill();
352 let _ = sh.wait();
353 let orphans = stragglers(&root);
354 assert_eq!(
355 orphans.len(),
356 2,
357 "the sleeps outlive their parent: {orphans:?}"
358 );
359 assert!(
360 !orphans.contains(&bystander.id()),
361 "a process outside the root was selected"
362 );
363 sweep(&root, grace);
364 let left = stragglers(&root);
365 let spared = bystander.try_wait().is_ok_and(|s| s.is_none());
366 let _ = bystander.kill();
367 let _ = bystander.wait();
368 assert!(spared, "the sweep stopped a process outside the root");
369 let _ = std::fs::remove_dir_all(&root);
370 assert!(left.is_empty(), "processes survived the sweep: {left:?}");
371 }
372
373 #[cfg(unix)]
374 #[test]
375 fn a_crashed_childs_orphans_are_swept() {
376 orphans_are_swept(
377 "term",
378 "sleep 300 & sleep 300 & wait",
379 Duration::from_secs(10),
380 );
381 }
382
383 #[cfg(unix)]
385 #[test]
386 fn orphans_that_ignore_sigterm_are_killed_after_the_grace_period() {
387 orphans_are_swept(
388 "kill",
389 "trap '' TERM; sleep 300 & sleep 300 & wait",
390 Duration::from_millis(300),
391 );
392 }
393}