Skip to main content

bynk/
sweep.rs

1//! Stopping what `bynk dev`'s wranglers leave behind (#1742).
2//!
3//! `bynk dev` stops each `wrangler dev` by signalling the child it spawned.
4//! That child is rarely the server. Resolved via npx, it is `npx`, and the
5//! server is four processes further down (`sh -c wrangler` → wrangler's `node`
6//! launcher → its CLI → two `workerd serve`s); a signal to `npx` reaches none of
7//! them. Resolved on PATH, the child is the launcher, which does pass SIGTERM
8//! on. But when a context's wrangler *crashes*, its `workerd`s are re-parented
9//! to init before `bynk dev` ever acts, and nothing walks from a child to them
10//! any more. Either way the survivors keep their ports, and the next `bynk dev`
11//! fails at boot with `bind(): Address already in use`.
12//!
13//! Two properties of those survivors outlast the re-parenting:
14//!
15//! - **Process group.** Nothing in the chain leaves the group it was spawned
16//!   into, which is `bynk dev`'s own: they share it so a Ctrl-C reaches them
17//!   all.
18//! - **Working directory.** `bynk dev` starts each wrangler inside its worker
19//!   dir under the managed build dir, and every descendant inherits it.
20//!
21//! [`sweep`] stops the processes that have both, so it reaches orphans without
22//! guessing by name. The working directory is what keeps it from touching an
23//! unrelated process: when `bynk dev` is not a group leader (a script without
24//! job control runs it), its group can hold the script's other jobs. Neither
25//! `bynk dev` itself nor its ancestors are ever selected.
26//!
27//! Unix only. Windows has neither process groups nor a working directory to
28//! read. There, `bynk dev` puts each wrangler in a job object instead, which
29//! holds the whole tree, orphans included, and stops it as one (#1762).
30
31use std::path::{Path, PathBuf};
32use std::time::Duration;
33
34/// Stop every process in this process group whose working directory is under
35/// `root`: SIGTERM first, then SIGKILL for what is still running after
36/// `grace`. Returns once none is left, or after SIGKILL has been given a moment
37/// to land.
38///
39/// Best effort, but never silently. If there is something to stop, it says
40/// so, because the wait can be long and would otherwise read as a hang. If
41/// anything is still running at the end, or `kill(1)` cannot be run, it names
42/// the survivors: they hold ports, and the next `bynk dev` would otherwise fail
43/// with a bind error and nothing pointing back here. If `ps` cannot be run
44/// there is nothing to select from, and this does nothing, as `bynk dev` did
45/// before #1742.
46pub 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                // Without `kill(1)` nothing can be signalled (std signals only
63                // its own children), so waiting out the grace would be wasted.
64                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                // SIGKILL is delivered asynchronously. Wait briefly for it to
72                // land, so a caller checking straight afterwards sees it.
73                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    // Windows: each wrangler's job object has already stopped its tree (#1762).
88    #[cfg(not(unix))]
89    let _ = (root, grace);
90}
91
92/// Name the processes [`sweep`] could not stop, and why.
93#[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/// The processes [`sweep`] would stop now.
107#[cfg(unix)]
108pub fn stragglers(root: &Path) -> Vec<u32> {
109    let Some(table) = process_table() else {
110        return Vec::new();
111    };
112    // Compare canonical paths: a process's cwd is reported resolved, and the
113    // build dir may sit behind a symlink (macOS's `/tmp` is `/private/tmp`).
114    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/// One row of `ps`: a process, its parent and its process group.
127#[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/// Every process on the system, from `ps` (POSIX `-A -o`, so Linux and macOS
136/// alike). `None` when `ps` cannot be run.
137#[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/// The other members of `me`'s process group, leaving out `me` and its
165/// ancestors. An ancestor shares the group when nothing between it and `me`
166/// started a new one, and stopping the shell or script that ran `bynk dev`
167/// would be the worst possible outcome.
168#[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    // Bounded, so a malformed table with a cycle cannot loop forever.
176    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/// The working directory of each of `pids` that still exists and can be read.
194#[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/// The working directory of each of `pids`, from one `lsof` call: there is no
202/// `/proc` here (macOS, the BSDs). Empty when `lsof` cannot be run.
203#[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    // `-a` ANDs the selections: only these pids, and only their cwd. `lsof`
214    // exits 1 when some pid has gone, so its status says nothing; parse what
215    // it printed.
216    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/// Parse `lsof -F pn` output: a `p<pid>` line opens each process, and each
227/// `n<path>` line after it names one of its files (here, only its cwd). Other
228/// field lines (`lsof` adds `f<fd>`) are skipped.
229#[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/// Send `signal` to `pids` with `kill(1)`; std can send only SIGKILL, and only
244/// to its own children.
245///
246/// `false` when `kill` could not be run at all. Its exit status is no test: a
247/// pid that has exited in the meantime makes it fail as well.
248#[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    /// `bynk dev` (300) run by a script (200) that has no job control, so all
276    /// three share the script's group with a sibling job (301). The orphan
277    /// (400, re-parented to init) is still in the group.
278    #[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),   // the login shell, its own group
283            row(200, 100, 200), // the script
284            row(300, 200, 200), // bynk dev
285            row(301, 200, 200), // another job of the script
286            row(310, 300, 200), // a wrangler bynk dev spawned
287            row(400, 1, 200),   // an orphaned workerd
288            row(500, 1, 500),   // unrelated
289        ];
290        let mut got = in_own_group(&table, 300);
291        got.sort_unstable();
292        // 301 is in the group too: the cwd filter, not the group, is what
293        // leaves it alone.
294        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    /// The #1742 regression, with `sh` in place of wrangler. Killing the `sh`
315    /// outright orphans its `sleep`s, the way a crashed wrangler orphans its
316    /// `workerd`s. Nothing reaches them from the child any more, and the sweep
317    /// must still stop them.
318    #[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        // A bystander: in the same process group (the test's), but working
342        // elsewhere, like a script's other job. The sweep must not touch it.
343        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    /// Orphans that ignore SIGTERM are killed once the grace period is over.
384    #[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}