From 1a2e833974c9b323f32ede8590ef447ecd328775 Mon Sep 17 00:00:00 2001 From: can1357 Date: Thu, 5 Feb 2026 08:10:52 +0100 Subject: [PATCH] feat: added Windows process handle duplication and signal handling for safe job termination - Added Windows-specific process handle duplication mechanism to safely manage process termination across multiple handles. - Implemented Windows signal handling support with Signal enum variants (Terminate, Kill, Interrupt) and kill_process() function using Windows API. - Added background job termination on cancellation with platform-specific implementations (TERM/KILL signals on Unix, TerminateProcess on Windows). - Implemented cancellation token support in shell execution with proper cleanup of background jobs when timeout or abort occurs. - Added test coverage for background job termination on timeout and abort scenarios. --- crates/brush-core-vendored/src/jobs.rs | 17 ++++ crates/brush-core-vendored/src/processes.rs | 45 +++++++++++ .../src/sys/stubs/signal.rs | 54 ++++++++++++- crates/pi-natives/src/shell.rs | 81 +++++++++++++++++-- .../coding-agent/test/bash-executor.test.ts | 40 +++++++++ 5 files changed, 230 insertions(+), 7 deletions(-) diff --git a/crates/brush-core-vendored/src/jobs.rs b/crates/brush-core-vendored/src/jobs.rs index 1b9c3aa0f..066b8e199 100644 --- a/crates/brush-core-vendored/src/jobs.rs +++ b/crates/brush-core-vendored/src/jobs.rs @@ -12,6 +12,9 @@ use crate::sys; use crate::trace_categories; use crate::traps; +#[cfg(windows)] +use std::os::windows::io::OwnedHandle; + pub(crate) type JobJoinHandle = tokio::task::JoinHandle>; pub(crate) type JobResult = (Job, Result); @@ -446,4 +449,18 @@ impl Job { // TODO: Don't assume that the first PID is the PGID. self.pgid.or_else(|| self.representative_pid()) } + + /// Duplicates process handles for termination on Windows. + #[cfg(windows)] + pub fn duplicate_kill_handles(&self) -> Vec { + let mut handles = Vec::new(); + for task in &self.tasks { + if let JobTask::External(process) = task { + if let Some(handle) = process.duplicate_kill_handle() { + handles.push(handle); + } + } + } + handles + } } diff --git a/crates/brush-core-vendored/src/processes.rs b/crates/brush-core-vendored/src/processes.rs index cf2cf3793..cd3e6fa27 100644 --- a/crates/brush-core-vendored/src/processes.rs +++ b/crates/brush-core-vendored/src/processes.rs @@ -3,6 +3,9 @@ use futures::FutureExt; use tokio_util::sync::CancellationToken; +#[cfg(windows)] +use std::os::windows::io::{AsRawHandle, OwnedHandle, RawHandle}; + use crate::{error, sys}; /// A waitable future that will yield the results of a child process's execution. @@ -16,14 +19,22 @@ pub struct ChildProcess { pid: Option, /// A waitable future that will yield the results of a child process's execution. exec_future: WaitableChildProcess, + #[cfg(windows)] + /// Windows handle duplicated from the child process for safe termination. + kill_handle: Option, } impl ChildProcess { /// Wraps a child process and its future. pub fn new(pid: Option, child: sys::process::Child) -> Self { + #[cfg(windows)] + let kill_handle = duplicate_handle(child.as_raw_handle()); + Self { pid, exec_future: Box::pin(child.wait_with_output()), + #[cfg(windows)] + kill_handle, } } @@ -32,6 +43,13 @@ impl ChildProcess { self.pid } + /// Duplicates the process handle for termination use on Windows. + #[cfg(windows)] + pub fn duplicate_kill_handle(&self) -> Option { + let handle = self.kill_handle.as_ref()?; + duplicate_handle(handle.as_raw_handle()) + } + /// Waits for the process to exit. /// /// If a cancellation token is provided and triggered, the process will be killed. @@ -113,6 +131,33 @@ impl ChildProcess { } } +#[cfg(windows)] +fn duplicate_handle(handle: RawHandle) -> Option { + use std::os::windows::io::FromRawHandle; + use windows_sys::Win32::System::Threading::{ + DuplicateHandle, GetCurrentProcess, DUPLICATE_SAME_ACCESS, + }; + + let current = unsafe { GetCurrentProcess() }; + let mut out_handle = std::ptr::null_mut(); + let ok = unsafe { + DuplicateHandle( + current, + handle as _, + current, + &mut out_handle, + 0, + 0, + DUPLICATE_SAME_ACCESS, + ) + }; + if ok == 0 || out_handle.is_null() { + return None; + } + + Some(unsafe { OwnedHandle::from_raw_handle(out_handle) }) +} + /// Represents the result of waiting for an executing process. pub enum ProcessWaitResult { /// The process completed. diff --git a/crates/brush-core-vendored/src/sys/stubs/signal.rs b/crates/brush-core-vendored/src/sys/stubs/signal.rs index 892e45294..4dff385c4 100644 --- a/crates/brush-core-vendored/src/sys/stubs/signal.rs +++ b/crates/brush-core-vendored/src/sys/stubs/signal.rs @@ -3,23 +3,55 @@ use crate::{error, sys, traps}; /// A stub enum representing system signals on unsupported platforms. +#[cfg(not(windows))] #[allow(unnameable_types)] #[derive(Clone, Copy, Eq, Hash, PartialEq)] pub enum Signal {} +/// Minimal signal representation for Windows. +#[cfg(windows)] +#[derive(Clone, Copy, Eq, Hash, PartialEq)] +pub enum Signal { + Terminate, + Kill, + Interrupt, +} + impl Signal { /// Returns an iterator over all possible signals. pub fn iterator() -> impl Iterator { - std::iter::empty() + #[cfg(windows)] + return [Self::Terminate, Self::Kill, Self::Interrupt].into_iter(); + #[cfg(not(windows))] + return std::iter::empty(); } /// Converts the signal into its corresponding name as a `&'static str`. pub const fn as_str(self) -> &'static str { + #[cfg(windows)] + { + return match self { + Self::Terminate => "TERM", + Self::Kill => "KILL", + Self::Interrupt => "INT", + }; + } + #[cfg(not(windows))] "" } /// Creates a `Signal` from a string representation. pub fn from_str(s: &str) -> Result { + #[cfg(windows)] + { + return match s.to_ascii_uppercase().as_str() { + "TERM" | "SIGTERM" => Ok(Self::Terminate), + "KILL" | "SIGKILL" => Ok(Self::Kill), + "INT" | "SIGINT" => Ok(Self::Interrupt), + _ => Err(error::ErrorKind::InvalidSignal(s.into()).into()), + }; + } + #[cfg(not(windows))] Err(error::ErrorKind::InvalidSignal(s.into()).into()) } } @@ -43,6 +75,26 @@ pub fn kill_process( _pid: sys::process::ProcessId, _signal: traps::TrapSignal, ) -> Result<(), error::Error> { + #[cfg(windows)] + { + use windows_sys::Win32::Foundation::CloseHandle; + use windows_sys::Win32::System::Threading::{OpenProcess, TerminateProcess, PROCESS_TERMINATE}; + + let pid = _pid as u32; + unsafe { + let handle = OpenProcess(PROCESS_TERMINATE, 0, pid); + if handle == 0 { + return Err(error::ErrorKind::FailedToSendSignal.into()); + } + let ok = TerminateProcess(handle, 1); + let _ = CloseHandle(handle); + if ok == 0 { + return Err(error::ErrorKind::FailedToSendSignal.into()); + } + } + return Ok(()); + } + #[cfg(not(windows))] Err(error::ErrorKind::NotSupportedOnThisPlatform("killing process").into()) } diff --git a/crates/pi-natives/src/shell.rs b/crates/pi-natives/src/shell.rs index ab37a2500..a1d0f77d3 100644 --- a/crates/pi-natives/src/shell.rs +++ b/crates/pi-natives/src/shell.rs @@ -12,8 +12,6 @@ //! }); //! ``` -#[cfg(windows)] -use std::collections::HashSet; use std::{ collections::HashMap, fs, @@ -22,6 +20,8 @@ use std::{ sync::Arc, time::Duration, }; +#[cfg(windows)] +use std::{collections::HashSet, os::windows::io::AsRawHandle}; #[cfg(windows)] mod windows; @@ -32,6 +32,7 @@ use brush_core::{ ProcessGroupPolicy, Shell as BrushShell, ShellValue, ShellVariable, builtins, env::EnvironmentScope, openfiles::{self, OpenFile, OpenFiles}, + sys, traps, }; use clap::Parser; use napi::{ @@ -44,6 +45,8 @@ use tokio::io::AsyncReadExt as _; use tokio_util::sync::CancellationToken; #[cfg(windows)] use windows::configure_windows_path; +#[cfg(windows)] +use windows_sys::Win32::System::Threading::TerminateProcess; use crate::task; @@ -524,7 +527,7 @@ async fn run_shell_command( params.set_fd(OpenFiles::STDOUT_FD, stdout_file); params.set_fd(OpenFiles::STDERR_FD, stderr_file); params.process_group_policy = ProcessGroupPolicy::NewProcessGroup; - params.set_cancel_token(cancel_token); + params.set_cancel_token(cancel_token.clone()); let mut env_scope_pushed = false; if let Some(env) = options.env.as_ref() { @@ -548,8 +551,9 @@ async fn run_shell_command( } } + let reader_cancel = cancel_token.clone(); let reader_handle = tokio::spawn(async move { - read_output(reader_file, on_chunk).await; + read_output(reader_file, on_chunk, reader_cancel).await; Result::<()>::Ok(()) }); let result = session @@ -557,6 +561,10 @@ async fn run_shell_command( .run_string(options.command.clone(), ¶ms) .await; + if cancel_token.is_cancelled() { + terminate_background_jobs(&session.shell); + } + if env_scope_pushed { session .shell @@ -572,6 +580,58 @@ async fn run_shell_command( result.map_err(|err| Error::from_reason(format!("Shell execution failed: {err}"))) } +#[cfg(unix)] +fn terminate_background_jobs(shell: &BrushShell) { + if shell.jobs.jobs.is_empty() { + return; + } + let Ok(signal) = "TERM".parse::() else { + return; + }; + let mut pgids = Vec::new(); + for job in &shell.jobs.jobs { + if let Some(pid) = job.process_group_id().or_else(|| job.representative_pid()) { + let _ = sys::signal::kill_process(pid, signal); + pgids.push(pid); + } + } + if pgids.is_empty() { + return; + } + tokio::spawn(async move { + time::sleep(Duration::from_millis(500)).await; + let Ok(signal) = "KILL".parse::() else { + return; + }; + for pid in pgids { + let _ = sys::signal::kill_process(pid, signal); + } + }); +} + +#[cfg(windows)] +fn terminate_background_jobs(shell: &BrushShell) { + if shell.jobs.jobs.is_empty() { + return; + } + let mut handles = Vec::new(); + for job in &shell.jobs.jobs { + handles.extend(job.duplicate_kill_handles()); + } + if handles.is_empty() { + return; + } + tokio::spawn(async move { + time::sleep(Duration::from_millis(500)).await; + for handle in handles { + // SAFETY: OwnedHandle keeps the duplicated handle alive for the duration. + unsafe { + let _ = TerminateProcess(handle.as_raw_handle() as _, 1); + } + } + }); +} + fn should_skip_env_var(key: &str) -> bool { if key.starts_with("BASH_FUNC_") && key.ends_with("%%") { return true; @@ -640,7 +700,11 @@ const fn session_keepalive(result: &ExecutionResult) -> bool { } } -async fn read_output(reader: fs::File, on_chunk: Option>) { +async fn read_output( + reader: fs::File, + on_chunk: Option>, + cancel_token: CancellationToken, +) { const REPLACEMENT: &str = "\u{FFFD}"; const BUF: usize = 4096; let mut buf = [0u8; BUF + 4]; // +4 for max UTF-8 char @@ -650,7 +714,12 @@ async fn read_output(reader: fs::File, on_chunk: Option res, + () = cancel_token.cancelled() => break, + } { Ok(0) => break, // EOF Ok(n) => n, Err(e) if e.kind() == io::ErrorKind::Interrupted => continue, diff --git a/packages/coding-agent/test/bash-executor.test.ts b/packages/coding-agent/test/bash-executor.test.ts index 871cd01f2..f717b8f2b 100644 --- a/packages/coding-agent/test/bash-executor.test.ts +++ b/packages/coding-agent/test/bash-executor.test.ts @@ -212,6 +212,46 @@ describe("executeBash", () => { expect(fs.existsSync(marker)).toBe(false); }); + it("kills background jobs on timeout", async () => { + if (process.platform === "win32") return; + + const marker = path.join(tempDir, "marker-bg.txt"); + const markerEscaped = marker.replace(/'/g, "'\\''"); + + const result = await executeBash(`{ sleep 2; echo done > '${markerEscaped}'; } & sleep 10`, { + cwd: tempDir, + timeout: 100, + }); + + expect(result.cancelled).toBe(true); + + await Bun.sleep(3000); + expect(fs.existsSync(marker)).toBe(false); + }); + + it("kills background jobs on abort", async () => { + if (process.platform === "win32") return; + + const marker = path.join(tempDir, "marker-bg-abort.txt"); + const markerEscaped = marker.replace(/'/g, "'\\''"); + const controller = new AbortController(); + + const promise = executeBash(`{ sleep 2; echo done > '${markerEscaped}'; } & sleep 10`, { + cwd: tempDir, + timeout: 10000, + signal: controller.signal, + }); + + await Bun.sleep(100); + controller.abort(); + const result = await promise; + + expect(result.cancelled).toBe(true); + + await Bun.sleep(3000); + expect(fs.existsSync(marker)).toBe(false); + }); + it("kills spawned process on abort (not just orphans it)", async () => { if (process.platform === "win32") return;