use std::path::{Path, PathBuf};
use std::process::{Command, Stdio};
use std::sync::{Arc, RwLock};
use super::WRAPPER_VARS;
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct ProcInfo {
pub(crate) exe: Option<PathBuf>,
pub(crate) args: Vec<String>,
pub(crate) env: Option<Vec<(String, String)>>,
pub(crate) start: Option<u64>,
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub(crate) struct Sockets {
pub(crate) listener: Option<u32>,
pub(crate) clients: Vec<u32>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum Signal {
Term,
Kill,
}
pub(crate) trait ProcessTable: Send + Sync {
fn listener(&self, port: u16) -> Result<Option<u32>, String>;
fn clients(&self, port: u16) -> Result<Vec<u32>, String>;
fn sockets(&self, port: u16) -> (Result<Option<u32>, String>, Result<Vec<u32>, String>) {
(self.listener(port), self.clients(port))
}
fn inspect(&self, pid: u32) -> Option<ProcInfo>;
fn spawn_detached(
&self,
exe: &Path,
args: &[String],
cwd: &Path,
env: &[(String, String)],
) -> Result<u32, String>;
fn signal(&self, pid: u32, sig: Signal) -> Result<(), String>;
fn is_alive(&self, pid: u32) -> bool;
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct ProcIdentity {
exe: Option<PathBuf>,
start: Option<u64>,
}
impl ProcIdentity {
pub(crate) fn of(info: &ProcInfo) -> Self {
Self {
exe: info.exe.clone(),
start: info.start,
}
}
fn is(&self, now: Option<&ProcInfo>) -> bool {
now.is_some_and(|n| n.exe == self.exe && n.start == self.start)
}
}
pub(crate) struct SystemTable;
const TERM_GRACE_POLLS: u32 = 30; const KILL_GRACE_POLLS: u32 = 20; const TERM_POLL: std::time::Duration = std::time::Duration::from_millis(100);
const TOOL_DEADLINE: std::time::Duration = std::time::Duration::from_secs(10);
pub(crate) fn poll_until<T>(
times: u32,
interval: std::time::Duration,
mut check: impl FnMut() -> Option<T>,
) -> Option<T> {
for i in 0..times {
if let Some(v) = check() {
return Some(v);
}
if i + 1 < times {
std::thread::sleep(interval);
}
}
None
}
pub(crate) fn terminate(
t: &dyn ProcessTable,
pid: u32,
expected: Option<&ProcIdentity>,
) -> Result<(), String> {
terminate_with(t, pid, expected, TERM_POLL)
}
fn terminate_with(
t: &dyn ProcessTable,
pid: u32,
expected: Option<&ProcIdentity>,
poll: std::time::Duration,
) -> Result<(), String> {
if pid == 0 {
return Err("refusing to signal pid 0".into());
}
let still_it = || expected.is_none_or(|id| id.is(t.inspect(pid).as_ref()));
let gone =
|polls: u32| poll_until(polls + 1, poll, || (!t.is_alive(pid)).then_some(())).is_some();
if !still_it() {
return Ok(());
}
t.signal(pid, Signal::Term)?;
if gone(TERM_GRACE_POLLS) || !still_it() {
return Ok(());
}
t.signal(pid, Signal::Kill)?;
if gone(KILL_GRACE_POLLS) {
return Ok(());
}
Err(format!("pid {pid} survived SIGTERM and SIGKILL"))
}
impl ProcessTable for SystemTable {
fn listener(&self, port: u16) -> Result<Option<u32>, String> {
Ok(sys::sockets(port)?.listener)
}
fn clients(&self, port: u16) -> Result<Vec<u32>, String> {
Ok(sys::sockets(port)?.clients)
}
fn sockets(&self, port: u16) -> (Result<Option<u32>, String>, Result<Vec<u32>, String>) {
match sys::sockets(port) {
Ok(s) => (Ok(s.listener), Ok(s.clients)),
Err(e) => (Err(e.clone()), Err(e)),
}
}
fn inspect(&self, pid: u32) -> Option<ProcInfo> {
if pid == 0 {
return None;
}
sys::inspect(pid)
}
fn spawn_detached(
&self,
exe: &Path,
args: &[String],
cwd: &Path,
env: &[(String, String)],
) -> Result<u32, String> {
let mut cmd = hub_command(exe, args, cwd, env);
#[cfg(unix)]
{
use std::os::unix::process::CommandExt;
unsafe {
cmd.pre_exec(|| {
libc::setsid();
Ok(())
});
}
}
#[cfg(windows)]
{
use std::os::windows::process::CommandExt;
const CREATE_NO_WINDOW: u32 = 0x0800_0000;
const CREATE_NEW_PROCESS_GROUP: u32 = 0x0000_0200;
cmd.creation_flags(CREATE_NO_WINDOW | CREATE_NEW_PROCESS_GROUP);
}
let mut child = cmd
.spawn()
.map_err(|e| format!("spawning {}: {e}", exe.display()))?;
let pid = child.id();
std::thread::Builder::new()
.name("cline-hub-reaper".into())
.spawn(move || {
let _ = child.wait();
})
.map_err(|e| format!("starting the reaper for pid {pid}: {e}"))?;
Ok(pid)
}
fn signal(&self, pid: u32, sig: Signal) -> Result<(), String> {
if pid == 0 {
return Err("refusing to signal pid 0".into());
}
match sig {
Signal::Term => sys::signal_term(pid),
Signal::Kill => sys::signal_kill(pid),
}
}
fn is_alive(&self, pid: u32) -> bool {
sys::is_alive(pid)
}
}
#[cfg_attr(not(any(target_os = "macos", windows)), allow(dead_code))]
fn run_tool(program: &str, args: &[&str]) -> crate::model_relay::trust_store::Ran {
use crate::model_relay::trust_store::{Runner, SystemRunner};
let args: Vec<String> = args.iter().map(|a| (*a).to_string()).collect();
SystemRunner.run(program, &args, TOOL_DEADLINE)
}
pub(crate) fn hub_command(
exe: &Path,
args: &[String],
cwd: &Path,
env: &[(String, String)],
) -> Command {
let mut c = Command::new(exe);
c.args(args)
.current_dir(cwd)
.stdin(Stdio::null())
.stdout(Stdio::null())
.stderr(Stdio::null());
for v in WRAPPER_VARS {
c.env_remove(v);
}
for (k, v) in env {
c.env(k, v);
}
c
}
static TABLE: RwLock<Option<Arc<dyn ProcessTable>>> = RwLock::new(None);
pub(crate) fn table() -> Arc<dyn ProcessTable> {
if let Some(t) = TABLE.read().unwrap_or_else(|e| e.into_inner()).clone() {
return t;
}
if cfg!(test) {
panic!(
"a test reached the real process table: install \
cline_app::process::install_for_tests(FakeTable) under HUB_SEAM_LOCK first"
);
}
Arc::new(SystemTable)
}
#[cfg(test)]
pub(crate) static HUB_SEAM_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
#[cfg(test)]
pub(crate) struct TableGuard {
_lock: std::sync::MutexGuard<'static, ()>,
}
#[cfg(test)]
impl Drop for TableGuard {
fn drop(&mut self) {
*TABLE.write().unwrap_or_else(|e| e.into_inner()) = None;
}
}
#[cfg(test)]
#[must_use = "the fake is uninstalled the moment this guard drops"]
pub(crate) fn install_for_tests(fake: Arc<dyn ProcessTable>) -> TableGuard {
let lock = HUB_SEAM_LOCK.lock().unwrap_or_else(|e| e.into_inner());
*TABLE.write().unwrap_or_else(|e| e.into_inner()) = Some(fake);
TableGuard { _lock: lock }
}
#[cfg(test)]
pub(crate) mod test_support {
use std::collections::{HashMap, HashSet, VecDeque};
use std::path::{Path, PathBuf};
use std::sync::Mutex;
use super::{ProcInfo, ProcessTable, Signal};
pub(crate) type Spawn = (PathBuf, Vec<String>, PathBuf, Vec<(String, String)>);
pub(crate) struct FakeTable {
pub(crate) expect_port: u16,
pub(crate) listener: Mutex<Result<Option<u32>, String>>,
pub(crate) clients: Mutex<Result<Vec<u32>, String>>,
pub(crate) procs: Mutex<HashMap<u32, ProcInfo>>,
pub(crate) spawns: Mutex<Vec<Spawn>>,
pub(crate) terminated: Mutex<Vec<u32>>,
pub(crate) killed: Mutex<Vec<u32>>,
pub(crate) on_spawn_listen: Mutex<Option<u32>>,
pub(crate) next_pid: Mutex<u32>,
pub(crate) inspect_script: Mutex<HashMap<u32, VecDeque<Option<ProcInfo>>>>,
pub(crate) survives_term: Mutex<bool>,
pub(crate) dead: Mutex<HashSet<u32>>,
}
impl FakeTable {
pub(crate) fn new(expect_port: u16) -> Self {
Self {
expect_port,
listener: Mutex::new(Ok(None)),
clients: Mutex::new(Ok(Vec::new())),
procs: Mutex::new(HashMap::new()),
spawns: Mutex::new(Vec::new()),
terminated: Mutex::new(Vec::new()),
killed: Mutex::new(Vec::new()),
on_spawn_listen: Mutex::new(None),
next_pid: Mutex::new(900),
inspect_script: Mutex::new(HashMap::new()),
survives_term: Mutex::new(false),
dead: Mutex::new(HashSet::new()),
}
}
fn check(&self, port: u16) {
assert_eq!(
port, self.expect_port,
"the fake process table was asked about port {port}, not {}",
self.expect_port
);
}
}
fn lock<T>(m: &Mutex<T>) -> std::sync::MutexGuard<'_, T> {
m.lock().unwrap_or_else(|e| e.into_inner())
}
impl ProcessTable for FakeTable {
fn listener(&self, port: u16) -> Result<Option<u32>, String> {
self.check(port);
lock(&self.listener).clone()
}
fn clients(&self, port: u16) -> Result<Vec<u32>, String> {
self.check(port);
lock(&self.clients).clone()
}
fn inspect(&self, pid: u32) -> Option<ProcInfo> {
if let Some(q) = lock(&self.inspect_script).get_mut(&pid) {
if q.len() > 1 {
return q.pop_front().flatten();
}
if let Some(last) = q.front() {
return last.clone();
}
}
lock(&self.procs).get(&pid).cloned()
}
fn spawn_detached(
&self,
exe: &Path,
args: &[String],
cwd: &Path,
env: &[(String, String)],
) -> Result<u32, String> {
lock(&self.spawns).push((
exe.to_path_buf(),
args.to_vec(),
cwd.to_path_buf(),
env.to_vec(),
));
let pid = {
let mut next = lock(&self.next_pid);
let pid = *next;
*next += 1;
pid
};
if let Some(listen) = lock(&self.on_spawn_listen).take() {
*lock(&self.listener) = Ok(Some(listen));
}
Ok(pid)
}
fn signal(&self, pid: u32, sig: Signal) -> Result<(), String> {
let ends = match sig {
Signal::Term => {
lock(&self.terminated).push(pid);
!*lock(&self.survives_term)
}
Signal::Kill => {
lock(&self.killed).push(pid);
true
}
};
if ends {
lock(&self.dead).insert(pid);
let mut listener = lock(&self.listener);
if matches!(*listener, Ok(Some(p)) if p == pid) {
*listener = Ok(None);
}
}
Ok(())
}
fn is_alive(&self, pid: u32) -> bool {
!lock(&self.dead).contains(&pid)
}
}
pub(crate) struct PanickingTable;
impl ProcessTable for PanickingTable {
fn listener(&self, port: u16) -> Result<Option<u32>, String> {
panic!("PanickingTable::listener({port}) was reached")
}
fn clients(&self, port: u16) -> Result<Vec<u32>, String> {
panic!("PanickingTable::clients({port}) was reached")
}
fn inspect(&self, pid: u32) -> Option<ProcInfo> {
panic!("PanickingTable::inspect({pid}) was reached")
}
fn spawn_detached(
&self,
exe: &Path,
_args: &[String],
_cwd: &Path,
_env: &[(String, String)],
) -> Result<u32, String> {
panic!(
"PanickingTable::spawn_detached({}) was reached",
exe.display()
)
}
fn signal(&self, pid: u32, sig: Signal) -> Result<(), String> {
panic!("PanickingTable::signal({pid}, {sig:?}) was reached")
}
fn is_alive(&self, pid: u32) -> bool {
panic!("PanickingTable::is_alive({pid}) was reached")
}
}
}
#[cfg_attr(not(target_os = "macos"), allow(dead_code))]
pub(crate) fn parse_lsof_pids(out: &str) -> Vec<u32> {
let mut pids = Vec::new();
for line in out.lines() {
if let Some(pid) = line
.strip_prefix('p')
.and_then(|p| p.trim().parse::<u32>().ok())
{
if !pids.contains(&pid) {
pids.push(pid);
}
}
}
pids
}
#[cfg_attr(not(target_os = "macos"), allow(dead_code))]
pub(crate) fn parse_lsof_sockets(out: &str) -> Sockets {
let (mut listen, mut established) = (String::new(), String::new());
let mut pid_line = None;
for line in out.lines() {
if line.starts_with('p') {
pid_line = Some(line);
} else if let (Some(state), Some(p)) = (line.strip_prefix("TST="), pid_line) {
let into = match state.trim() {
"LISTEN" => &mut listen,
"ESTABLISHED" => &mut established,
_ => continue,
};
into.push_str(p);
into.push('\n');
}
}
Sockets {
listener: parse_lsof_pids(&listen).first().copied(),
clients: parse_lsof_pids(&established),
}
}
#[cfg_attr(not(target_os = "macos"), allow(dead_code))]
pub(crate) fn parse_procargs2(buf: &[u8]) -> Option<ProcInfo> {
let argc_bytes: [u8; 4] = buf.get(..4)?.try_into().ok()?;
let argc = usize::try_from(i32::from_ne_bytes(argc_bytes)).ok()?;
let mut pos = 4;
let exe = next_cstr(buf, &mut pos)?;
while buf.get(pos) == Some(&0) {
pos += 1;
}
let mut args = Vec::with_capacity(argc);
for _ in 0..argc {
args.push(next_cstr(buf, &mut pos)?);
}
let mut env = Vec::new();
while pos < buf.len() {
let entry = next_cstr(buf, &mut pos)?;
if entry.is_empty() {
break;
}
if let Some((k, v)) = entry.split_once('=') {
env.push((k.to_string(), v.to_string()));
}
}
Some(ProcInfo {
exe: (!exe.is_empty()).then(|| PathBuf::from(exe)),
args,
env: Some(env),
start: None,
})
}
fn next_cstr(buf: &[u8], pos: &mut usize) -> Option<String> {
let rest = buf.get(*pos..)?;
let end = rest.iter().position(|b| *b == 0)?;
let s = String::from_utf8_lossy(&rest[..end]).into_owned();
*pos += end + 1;
Some(s)
}
#[cfg_attr(not(target_os = "linux"), allow(dead_code))]
pub(crate) fn parse_proc_net_tcp(table: &str, port: u16, state: u8) -> Vec<u64> {
let port_of = |addr: &str| {
addr.rsplit_once(':')
.and_then(|(_, p)| u16::from_str_radix(p, 16).ok())
};
let mut inodes = Vec::new();
for line in table.lines() {
let f: Vec<&str> = line.split_whitespace().collect();
if f.len() < 10 {
continue;
}
let Ok(st) = u8::from_str_radix(f[3], 16) else {
continue;
};
if st != state {
continue;
}
let local = port_of(f[1]) == Some(port);
let remote = port_of(f[2]) == Some(port);
let hit = if state == 0x0A {
local
} else {
local || remote
};
if !hit {
continue;
}
if let Ok(inode) = f[9].parse::<u64>() {
if inode != 0 && !inodes.contains(&inode) {
inodes.push(inode);
}
}
}
inodes
}
#[cfg_attr(not(windows), allow(dead_code))]
pub(crate) fn parse_netstat_ano(out: &str, port: u16) -> (Option<u32>, Vec<u32>) {
let suffix = format!(":{port}");
let listen_addr = format!("127.0.0.1:{port}");
let mut listener = None;
let mut clients = Vec::new();
for line in out.lines() {
let f: Vec<&str> = line.split_whitespace().collect();
if f.len() < 5 || !f[0].eq_ignore_ascii_case("TCP") {
continue;
}
let Ok(pid) = f[4].parse::<u32>() else {
continue;
};
match f[3] {
"LISTENING" if f[1] == listen_addr && listener.is_none() => listener = Some(pid),
"ESTABLISHED"
if (f[1].ends_with(&suffix) || f[2].ends_with(&suffix))
&& !clients.contains(&pid) =>
{
clients.push(pid)
}
_ => {}
}
}
(listener, clients)
}
#[cfg_attr(not(windows), allow(dead_code))]
pub(crate) fn parse_cim_json(out: &str) -> Option<ProcInfo> {
let v: serde_json::Value = serde_json::from_str(out.trim()).ok()?;
let obj = v.as_object()?;
let exe = obj
.get("ExecutablePath")
.and_then(|e| e.as_str())
.filter(|e| !e.is_empty())
.map(PathBuf::from);
let args = obj
.get("CommandLine")
.and_then(|c| c.as_str())
.map(|c| vec![c.to_string()])
.unwrap_or_default();
let start = obj.get("Start").and_then(serde_json::Value::as_u64);
Some(ProcInfo {
exe,
args,
env: None,
start,
})
}
#[cfg_attr(not(target_os = "linux"), allow(dead_code))]
pub(crate) fn parse_proc_stat_start(stat: &str) -> Option<u64> {
let (_, rest) = stat.rsplit_once(')')?;
rest.split_whitespace().nth(19)?.parse().ok()
}
#[cfg(target_os = "macos")]
mod sys {
use super::{parse_lsof_sockets, parse_procargs2, run_tool, ProcInfo, Sockets};
use crate::model_relay::trust_store::Ran;
pub(super) fn sockets(port: u16) -> Result<Sockets, String> {
let spec = format!("-iTCP:{port}");
let lsof = "/usr/sbin/lsof"; match run_tool(lsof, &["-nP", &spec, "-FpT"]) {
Ran::Exited {
success,
code,
stdout,
stderr,
} => {
if success || (code == Some(1) && stdout.trim().is_empty()) {
Ok(parse_lsof_sockets(&stdout))
} else {
Err(format!("lsof exited {code:?}: {}", stderr.trim()))
}
}
Ran::Missing => Err("lsof: not found".into()),
Ran::TimedOut => Err("lsof: no answer within the deadline".into()),
Ran::Failed(e) => Err(format!("lsof: {e}")),
}
}
fn procargs2(pid: u32) -> Option<Vec<u8>> {
let pid = libc::c_int::try_from(pid).ok()?;
let mut argmax: libc::c_int = 0;
let mut size = std::mem::size_of::<libc::c_int>();
let mut mib = [libc::CTL_KERN, libc::KERN_ARGMAX];
let rc = unsafe {
libc::sysctl(
mib.as_mut_ptr(),
2,
(&mut argmax as *mut libc::c_int).cast(),
&mut size,
std::ptr::null_mut(),
0,
)
};
if rc != 0 || argmax <= 0 {
return None;
}
let mut buf = vec![0u8; usize::try_from(argmax).ok()?];
let mut len = buf.len();
let mut mib = [libc::CTL_KERN, libc::KERN_PROCARGS2, pid];
let rc = unsafe {
libc::sysctl(
mib.as_mut_ptr(),
3,
buf.as_mut_ptr().cast(),
&mut len,
std::ptr::null_mut(),
0,
)
};
if rc != 0 {
return None;
}
buf.truncate(len);
Some(buf)
}
fn start(pid: u32) -> Option<u64> {
let pid = libc::c_int::try_from(pid).ok()?;
let mut info: libc::proc_bsdinfo = unsafe { std::mem::zeroed() };
let size = libc::c_int::try_from(std::mem::size_of::<libc::proc_bsdinfo>()).ok()?;
let rc = unsafe {
libc::proc_pidinfo(
pid,
libc::PROC_PIDTBSDINFO,
0,
(&mut info as *mut libc::proc_bsdinfo).cast(),
size,
)
};
if rc != size {
return None;
}
info.pbi_start_tvsec
.checked_mul(1_000_000)?
.checked_add(info.pbi_start_tvusec)
}
pub(super) fn inspect(pid: u32) -> Option<ProcInfo> {
let mut info = parse_procargs2(&procargs2(pid)?)?;
info.start = start(pid);
Some(info)
}
pub(super) use super::unix_signals::{is_alive, signal_kill, signal_term};
}
#[cfg(target_os = "linux")]
mod sys {
use std::path::PathBuf;
use super::{parse_proc_net_tcp, parse_proc_stat_start, ProcInfo, Sockets};
const PROC: &str = "/proc";
fn inodes(port: u16, state: u8) -> Result<Vec<u64>, String> {
let mut all = Vec::new();
let mut read_any = false;
for table in ["tcp", "tcp6"] {
let path = PathBuf::from(PROC).join("net").join(table);
if let Ok(body) = std::fs::read_to_string(&path) {
read_any = true;
for inode in parse_proc_net_tcp(&body, port, state) {
if !all.contains(&inode) {
all.push(inode);
}
}
}
}
if read_any {
Ok(all)
} else {
Err("neither /proc/net/tcp nor /proc/net/tcp6 is readable".into())
}
}
fn holders(inodes: &[u64]) -> Result<Vec<u32>, String> {
if inodes.is_empty() {
return Ok(Vec::new());
}
let wanted: Vec<String> = inodes.iter().map(|i| format!("socket:[{i}]")).collect();
let entries = std::fs::read_dir(PROC).map_err(|e| format!("/proc: {e}"))?;
let mut pids = Vec::new();
for entry in entries.flatten() {
let Some(pid) = entry
.file_name()
.to_str()
.and_then(|n| n.parse::<u32>().ok())
else {
continue;
};
let Ok(fds) = std::fs::read_dir(entry.path().join("fd")) else {
continue;
};
let holds = fds.flatten().any(|fd| {
std::fs::read_link(fd.path())
.ok()
.is_some_and(|l| wanted.iter().any(|w| l.as_os_str() == w.as_str()))
});
if holds && !pids.contains(&pid) {
pids.push(pid);
}
}
Ok(pids)
}
pub(super) fn listener(port: u16) -> Result<Option<u32>, String> {
Ok(holders(&inodes(port, 0x0A)?)?.first().copied())
}
pub(super) fn clients(port: u16) -> Result<Vec<u32>, String> {
holders(&inodes(port, 0x01)?)
}
pub(super) fn sockets(port: u16) -> Result<Sockets, String> {
Ok(Sockets {
listener: listener(port)?,
clients: clients(port)?,
})
}
fn nul_split(raw: &[u8]) -> Vec<String> {
raw.split(|b| *b == 0)
.filter(|s| !s.is_empty())
.map(|s| String::from_utf8_lossy(s).into_owned())
.collect()
}
pub(super) fn inspect(pid: u32) -> Option<ProcInfo> {
let dir = PathBuf::from(PROC).join(pid.to_string());
let args = nul_split(&std::fs::read(dir.join("cmdline")).ok()?);
let exe = std::fs::read_link(dir.join("exe")).ok();
let env = std::fs::read(dir.join("environ")).ok().map(|raw| {
nul_split(&raw)
.into_iter()
.filter_map(|e| {
e.split_once('=')
.map(|(k, v)| (k.to_string(), v.to_string()))
})
.collect()
});
let start = std::fs::read_to_string(dir.join("stat"))
.ok()
.and_then(|s| parse_proc_stat_start(&s));
Some(ProcInfo {
exe,
args,
env,
start,
})
}
pub(super) use super::unix_signals::{is_alive, signal_kill, signal_term};
}
#[cfg(windows)]
mod sys {
use super::{parse_cim_json, parse_netstat_ano, run_tool, ProcInfo, Sockets};
use crate::model_relay::trust_store::Ran;
pub(super) fn sockets(port: u16) -> Result<Sockets, String> {
match run_tool("netstat", &["-ano", "-p", "TCP"]) {
Ran::Exited {
success: true,
stdout,
..
} => {
let (listener, clients) = parse_netstat_ano(&stdout, port);
Ok(Sockets { listener, clients })
}
Ran::Exited { code, .. } => Err(format!("netstat exited {code:?}")),
Ran::Missing => Err("netstat: not found".into()),
Ran::TimedOut => Err("netstat: no answer within the deadline".into()),
Ran::Failed(e) => Err(format!("netstat: {e}")),
}
}
pub(super) fn inspect(pid: u32) -> Option<ProcInfo> {
let script = format!(
"Get-CimInstance Win32_Process -Filter 'ProcessId={pid}' | Select-Object \
ExecutablePath,CommandLine,@{{n='Start';e={{$_.CreationDate.ToFileTimeUtc()}}}} \
| ConvertTo-Json"
);
match run_tool("powershell", &["-NoProfile", "-Command", &script]) {
Ran::Exited {
success: true,
stdout,
..
} => parse_cim_json(&stdout),
_ => None,
}
}
fn taskkill(pid: u32, force: bool) -> Result<(), String> {
let pid_s = pid.to_string();
let mut args = vec!["/PID", pid_s.as_str(), "/T"];
if force {
args.push("/F");
}
match run_tool("taskkill", &args) {
Ran::Exited { .. } => Ok(()),
Ran::Missing => Err(format!("taskkill {pid}: not found")),
Ran::TimedOut => Err(format!("taskkill {pid}: no answer within the deadline")),
Ran::Failed(e) => Err(format!("taskkill {pid}: {e}")),
}
}
pub(super) fn signal_term(pid: u32) -> Result<(), String> {
taskkill(pid, false)
}
pub(super) fn signal_kill(pid: u32) -> Result<(), String> {
taskkill(pid, true)
}
pub(super) fn is_alive(pid: u32) -> bool {
let filter = format!("PID eq {pid}");
match run_tool("tasklist", &["/FI", &filter, "/NH", "/FO", "CSV"]) {
Ran::Exited { stdout, .. } => stdout.contains(&format!("\"{pid}\"")),
_ => true,
}
}
}
#[cfg(not(any(target_os = "macos", target_os = "linux", windows)))]
mod sys {
use super::{ProcInfo, Sockets};
const UNSUPPORTED: &str = "process table not supported on this platform";
pub(super) fn sockets(_port: u16) -> Result<Sockets, String> {
Err(UNSUPPORTED.into())
}
pub(super) fn inspect(_pid: u32) -> Option<ProcInfo> {
None
}
#[cfg(unix)]
pub(super) use super::unix_signals::{is_alive, signal_kill, signal_term};
#[cfg(not(unix))]
pub(super) fn signal_term(_pid: u32) -> Result<(), String> {
Err(UNSUPPORTED.into())
}
#[cfg(not(unix))]
pub(super) fn signal_kill(_pid: u32) -> Result<(), String> {
Err(UNSUPPORTED.into())
}
#[cfg(not(unix))]
pub(super) fn is_alive(_pid: u32) -> bool {
true
}
}
#[cfg(unix)]
mod unix_signals {
fn to_pid(pid: u32) -> Result<libc::pid_t, String> {
libc::pid_t::try_from(pid)
.ok()
.filter(|p| *p > 0)
.ok_or_else(|| format!("refusing to signal pid {pid}"))
}
fn signal(pid: u32, sig: libc::c_int) -> Result<(), String> {
let p = to_pid(pid)?;
let rc = unsafe { libc::kill(p, sig) };
if rc == 0 {
return Ok(());
}
let err = std::io::Error::last_os_error();
if err.raw_os_error() == Some(libc::ESRCH) {
return Ok(()); }
Err(format!("signalling pid {pid}: {err}"))
}
pub(super) fn signal_term(pid: u32) -> Result<(), String> {
signal(pid, libc::SIGTERM)
}
pub(super) fn signal_kill(pid: u32) -> Result<(), String> {
signal(pid, libc::SIGKILL)
}
pub(super) fn is_alive(pid: u32) -> bool {
let Ok(p) = to_pid(pid) else { return false };
let rc = unsafe { libc::kill(p, 0) };
let exists = rc == 0 || std::io::Error::last_os_error().raw_os_error() == Some(libc::EPERM);
exists && !is_zombie(pid)
}
#[cfg(target_os = "linux")]
fn is_zombie(pid: u32) -> bool {
std::fs::read_to_string(
std::path::PathBuf::from("/proc")
.join(pid.to_string())
.join("stat"),
) .is_ok_and(|s| crate::core::proc_stat::proc_stat_says_exited(&s))
}
#[cfg(not(target_os = "linux"))]
fn is_zombie(_pid: u32) -> bool {
false
}
}
#[cfg(test)]
mod tests {
use super::*;
fn procargs_buffer() -> Vec<u8> {
let mut b = 3i32.to_ne_bytes().to_vec();
b.extend_from_slice(b"/x/code-sidecar\0\0\0\0");
b.extend_from_slice(b"code-sidecar\0--cline-hub-daemon\0--port\0");
b.extend_from_slice(b"HTTPS_PROXY=http://127.0.0.1:7600\0PATH=/usr/bin\0\0");
b
}
#[test]
fn procargs2_splits_exe_args_and_env() {
let info = parse_procargs2(&procargs_buffer()).expect("parses");
assert_eq!(info.exe, Some(PathBuf::from("/x/code-sidecar")));
assert_eq!(
info.args,
vec!["code-sidecar", "--cline-hub-daemon", "--port"]
);
assert_eq!(
info.env,
Some(vec![
(
"HTTPS_PROXY".to_string(),
"http://127.0.0.1:7600".to_string()
),
("PATH".to_string(), "/usr/bin".to_string()),
])
);
}
#[test]
fn procargs2_refuses_a_truncated_buffer() {
let full = procargs_buffer();
let cut = full
.windows(b"--cline".len())
.position(|w| w == b"--cline")
.expect("the second arg")
+ 4;
assert_eq!(parse_procargs2(&full[..cut]), None);
assert_eq!(parse_procargs2(&full[..2]), None);
}
#[test]
fn lsof_pids_parses_field_output() {
assert_eq!(
parse_lsof_pids("p31077\nf12\np31077\np88893\n"),
vec![31077, 88893]
);
}
const TCP: &str = " sl local_address rem_address st tx_queue rx_queue tr tm->when retrnsmt uid timeout inode\n\
0: 0100007F:6377 00000000:0000 0A 00000000:00000000 00:00000000 00000000 501 0 11111 1 0000000000000000 100 0 0 10 0\n\
1: 0100007F:6377 0100007F:C350 01 00000000:00000000 00:00000000 00000000 501 0 22222 1 0000000000000000 20 4 30 10 -1\n\
2: 0100007F:C350 0100007F:6377 01 00000000:00000000 00:00000000 00000000 501 0 33333 1 0000000000000000 20 4 30 10 -1\n\
3: 0100007F:1F90 0100007F:C351 01 00000000:00000000 00:00000000 00000000 501 0 44444 1 0000000000000000 20 4 30 10 -1\n";
#[test]
fn proc_net_tcp_finds_the_listener_inode() {
assert_eq!(parse_proc_net_tcp(TCP, 25463, 0x0A), vec![11111]);
let established_only = TCP
.lines()
.filter(|l| !l.contains(" 0A "))
.collect::<Vec<_>>()
.join("\n");
assert!(
parse_proc_net_tcp(&established_only, 25463, 0x0A).is_empty(),
"a 01 row on the same port is not a listener"
);
}
#[test]
fn proc_net_tcp_finds_client_inodes() {
assert_eq!(parse_proc_net_tcp(TCP, 25463, 0x01), vec![22222, 33333]);
assert!(
parse_proc_net_tcp(TCP, 25999, 0x01).is_empty(),
"a 01 row on another port is none of the hub's"
);
}
#[test]
fn netstat_ano_finds_listener_and_clients() {
let out = "\nActive Connections\n\n Proto Local Address Foreign Address State PID\n\
TCP 127.0.0.1:25463 0.0.0.0:0 LISTENING 31077\n\
TCP 127.0.0.1:25463 127.0.0.1:50000 ESTABLISHED 31077\n\
TCP 127.0.0.1:50000 127.0.0.1:25463 ESTABLISHED 5120\n\
TCP 127.0.0.1:8080 127.0.0.1:50001 ESTABLISHED 777\n";
assert_eq!(
parse_netstat_ano(out, 25463),
(Some(31077), vec![31077, 5120])
);
}
#[test]
fn cim_json_reads_exe_and_command_line() {
let with = r#"{"ExecutablePath":"C:\\x\\code-sidecar.exe","CommandLine":"C:\\x\\code-sidecar.exe --cline-hub-daemon --port 25463","Start":133712345678901234}"#;
let info = parse_cim_json(with).expect("parses");
assert_eq!(info.exe, Some(PathBuf::from("C:\\x\\code-sidecar.exe")));
assert_eq!(info.args.len(), 1);
assert!(info.args[0].contains(super::super::HUB_ARG));
assert_eq!(info.env, None);
assert_eq!(info.start, Some(133_712_345_678_901_234));
let without =
r#"{"ExecutablePath":"C:\\x\\code-sidecar.exe","CommandLine":"x --cline-hub-daemon"}"#;
assert_eq!(parse_cim_json(without).expect("parses").start, None);
}
#[test]
fn proc_stat_start_reads_field_22_after_the_last_paren() {
let stat = "1234 (code (side) car) S 1 1234 1234 0 -1 4194560 10 0 0 0 1 2 0 0 20 0 7 0 \
987654 1000 200 18446744073709551615";
assert_eq!(parse_proc_stat_start(stat), Some(987_654));
let cut = "1234 (code (side) car) S 1 1234 1234 0 -1 4194560 10 0 0 0 1 2 0 0 20 0 7 0";
assert_eq!(parse_proc_stat_start(cut), None);
}
#[test]
fn hub_command_strips_every_wrapper_variable() {
use std::collections::BTreeSet;
use std::ffi::{OsStr, OsString};
let key = |k: &str| -> OsString {
if cfg!(windows) {
k.to_ascii_uppercase().into()
} else {
k.into()
}
};
let envs = |c: &Command| -> Vec<(OsString, Option<OsString>)> {
c.get_envs()
.map(|(k, v)| (key(&k.to_string_lossy()), v.map(OsStr::to_os_string)))
.collect()
};
let expected = WRAPPER_VARS
.iter()
.map(|v| key(v))
.collect::<BTreeSet<_>>()
.len();
let exe = Path::new("/x/code-sidecar");
let args = vec!["--cline-hub-daemon".to_string()];
let cwd = Path::new("/");
let clean = envs(&hub_command(exe, &args, cwd, &[]));
assert_eq!(clean.len(), expected, "{clean:?}");
for v in WRAPPER_VARS {
assert!(clean.contains(&(key(v), None)), "{v} is removed: {clean:?}");
}
let ours = super::super::ProxyEnvVars {
proxy_url: "http://127.0.0.1:7600".to_string(),
ca_pem: PathBuf::from("/x/ca.pem"),
no_proxy: "127.0.0.1:7600".to_string(),
}
.overrides(None, None);
let with = envs(&hub_command(exe, &args, cwd, &ours));
for (k, v) in &ours {
assert!(
with.contains(&(key(k), Some(OsString::from(v)))),
"{k} is ours: {with:?}"
);
}
assert_eq!(with.len(), expected);
}
}
#[cfg(test)]
mod review_tests {
use std::collections::VecDeque;
use std::time::Duration;
use super::test_support::FakeTable;
use super::*;
fn hub() -> ProcInfo {
ProcInfo {
exe: Some(PathBuf::from("/x/code-sidecar")),
args: vec!["--cline-hub-daemon".to_string()],
env: None,
start: Some(1),
}
}
#[test]
fn a_pid_reused_during_the_term_grace_is_never_killed() {
let mut reused = hub();
reused.start = Some(2);
for (after_term, killed) in [(Some(reused), false), (None, false), (Some(hub()), true)] {
let t = FakeTable::new(1);
*t.survives_term.lock().unwrap() = true;
t.inspect_script
.lock()
.unwrap()
.insert(500, VecDeque::from([Some(hub()), after_term.clone()]));
let r = terminate_with(&t, 500, Some(&ProcIdentity::of(&hub())), Duration::ZERO);
assert!(r.is_ok(), "{after_term:?}: {r:?}");
assert_eq!(*t.terminated.lock().unwrap(), vec![500], "{after_term:?}");
assert_eq!(
!t.killed.lock().unwrap().is_empty(),
killed,
"{after_term:?}"
);
}
}
#[test]
fn an_unverified_child_is_signalled() {
let t = FakeTable::new(1);
terminate_with(&t, 900, None, Duration::ZERO).expect("gone");
assert_eq!(*t.terminated.lock().unwrap(), vec![900]);
assert!(t.killed.lock().unwrap().is_empty());
}
#[test]
fn poll_until_checks_first_and_never_sleeps_after_the_last() {
let mut calls = 0;
assert_eq!(
poll_until(3, Duration::ZERO, || {
calls += 1;
None::<()>
}),
None
);
assert_eq!(calls, 3);
let mut calls = 0;
assert_eq!(
poll_until(5, Duration::ZERO, || {
calls += 1;
(calls == 2).then_some(calls)
}),
Some(2)
);
assert_eq!(poll_until(0, Duration::ZERO, || Some(())), None);
}
#[test]
fn lsof_sockets_splits_listen_and_established() {
let out = "p26852\nf3\nTST=LISTEN\nTQR=0\nTQS=0\nf4\nTST=ESTABLISHED\nTQR=0\nTQS=0\n\
f5\nTST=ESTABLISHED\nTQR=0\nTQS=0\np777\nf9\nTST=ESTABLISHED\nf10\nTST=CLOSE_WAIT\n";
assert_eq!(
parse_lsof_sockets(out),
Sockets {
listener: Some(26852),
clients: vec![26852, 777],
}
);
assert_eq!(parse_lsof_sockets(""), Sockets::default());
}
}