From 3e38a5b514374774e5cee96b9129ccc757fd1ab0 Mon Sep 17 00:00:00 2001 From: roboomp Date: Thu, 16 Jul 2026 07:22:58 +0000 Subject: [PATCH 1/2] fix(bash): drained piped output before timeout return - Delayed reader cancellation so pipeline consumers can flush after producers are terminated. - Kept the JavaScript watchdog behind bounded native timeout cleanup. - Added native and executor regressions for timeout-time output draining. Fixes #5316 --- crates/pi-natives/src/shell.rs | 27 +++++++++++ crates/pi-shell/src/shell.rs | 5 ++ packages/coding-agent/CHANGELOG.md | 4 ++ .../coding-agent/src/exec/bash-executor.ts | 15 ++++-- .../coding-agent/test/bash-executor.test.ts | 48 ++++++++++++++++--- packages/natives/CHANGELOG.md | 4 ++ 6 files changed, 92 insertions(+), 11 deletions(-) diff --git a/crates/pi-natives/src/shell.rs b/crates/pi-natives/src/shell.rs index ccc0e07dd..4b80406f4 100644 --- a/crates/pi-natives/src/shell.rs +++ b/crates/pi-natives/src/shell.rs @@ -558,4 +558,31 @@ mod tests { .expect("shell run should return"); assert!(result.cancelled); } + + #[tokio::test(flavor = "multi_thread")] + async fn timeout_drains_pipeline_output_before_stopping_reader() { + let shell = CoreShell::new(None); + let (tx, rx) = flume::unbounded::(); + let result = shell + .run( + CoreShellRunOptions { + command: "yes x | tail -5".to_string(), + cwd: None, + env: None, + timeout_ms: Some(50), + }, + Some(tx), + CancelToken::new(Some(50)), + ) + .await + .expect("shell run"); + + let mut output = String::new(); + while let Ok(chunk) = rx.recv_async().await { + output.push_str(&chunk); + } + + assert!(result.timed_out); + assert_eq!(output.lines().filter(|line| *line == "x").count(), 5); + } } diff --git a/crates/pi-shell/src/shell.rs b/crates/pi-shell/src/shell.rs index f2254a258..369d61c13 100644 --- a/crates/pi-shell/src/shell.rs +++ b/crates/pi-shell/src/shell.rs @@ -1141,11 +1141,16 @@ async fn run_shell_command_once( } } }); + // Let pipeline consumers flush output after cancellation kills their + // producers. The outer run cancellation remains bounded, and this delayed + // fallback still releases readers whose writers never close. + const CANCEL_READER_GRACE: Duration = Duration::from_millis(500); let cancel_bridge = tokio::spawn({ let cancel_token = cancel_token.clone(); let reader_cancel = reader_cancel.clone(); async move { cancel_token.cancelled().await; + time::sleep(CANCEL_READER_GRACE).await; reader_cancel.cancel(); } }); diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index e0b1b1608..8bcefe0f7 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -10,6 +10,10 @@ - Updated status event log to prioritize the most recent entries in the display window +### Fixed + +- Fixed Windows bash crashes when a piped command times out while flushing output; explicit-timeout watchdogs now wait for bounded native teardown instead of returning mid-drain. ([#5316](https://github.com/can1357/oh-my-pi/issues/5316)) + ### Removed - Removed the unreliable Bing and Yahoo HTML-scraping web search providers diff --git a/packages/coding-agent/src/exec/bash-executor.ts b/packages/coding-agent/src/exec/bash-executor.ts index 505feca96..ef75a9256 100644 --- a/packages/coding-agent/src/exec/bash-executor.ts +++ b/packages/coding-agent/src/exec/bash-executor.ts @@ -70,6 +70,9 @@ const shellSessionsInUse = new Set(); */ const retainedShells = new Set(); const RETAIN_REAP_INTERVAL_MS = 5_000; +// Native cancellation may spend two seconds unwinding the shell before its +// N-API chunk bridge drains. The JS watchdog must not race that teardown. +const NATIVE_TIMEOUT_FALLBACK_GRACE_MS = 5_000; async function retainShellWithLiveBackgroundJobs(shell: Shell): Promise { let live: number; @@ -302,16 +305,18 @@ export async function executeBash(command: string, options?: BashExecutorOptions const nativeTimeoutMs = requestedTimeoutMs !== undefined && requestedTimeoutMs > 0 ? requestedTimeoutMs : undefined; const nativeOwnsTimeout = nativeTimeoutMs !== undefined; if (deadlineTimeoutMs !== undefined) { + const fallbackTimeoutMs = nativeOwnsTimeout + ? deadlineTimeoutMs + NATIVE_TIMEOUT_FALLBACK_GRACE_MS + : deadlineTimeoutMs; timeoutTimer = setTimeout(() => { - // Explicit timeouts are already enforced inside pi-natives via - // `timeoutMs`. Do not also abort the JS AbortSignal here: on Windows, - // aborting that signal while a piped command is still forwarding output - // can terminate the Bun host before the native timeout result resolves. + // Explicit timeouts are enforced inside pi-natives via `timeoutMs`. + // Give native cancellation time to flush pipeline output and drain the + // N-API bridge before this result-only watchdog quarantines the run. if (!nativeOwnsTimeout) { abortCurrentExecution(); } timeoutDeferred.resolve("timeout"); - }, deadlineTimeoutMs); + }, fallbackTimeoutMs); } let resetSession = false; diff --git a/packages/coding-agent/test/bash-executor.test.ts b/packages/coding-agent/test/bash-executor.test.ts index f428d3ea9..06ce9dbe1 100644 --- a/packages/coding-agent/test/bash-executor.test.ts +++ b/packages/coding-agent/test/bash-executor.test.ts @@ -520,19 +520,54 @@ exit 64 expect(next.output.trim()).toBe("still_persistent"); }); - it("does not abort the native signal when the JavaScript timeout fallback returns streamed output", async () => { - // Compress the JS-side fallback timer (floored at 1000ms in the source) so - // the safety-net fires deterministically without a real 1s wait. Only long - // timers are shrunk — fs/subprocess setup keeps real scheduling — and the - // reported "1 seconds" derives from the configured timeout, not the timer. + it("waits for native timeout teardown to flush piped output", async () => { const realSetTimeout = globalThis.setTimeout; vi.spyOn(globalThis, "setTimeout").mockImplementation(((handler: () => void, ms?: number, ...rest: unknown[]) => realSetTimeout( handler, - typeof ms === "number" && ms >= 1000 ? 5 : ms, + ms === 1000 ? 5 : typeof ms === "number" && ms > 1000 ? 50 : ms, ...rest, )) as typeof globalThis.setTimeout); + let nativeSignal: AbortSignal | undefined; + vi.spyOn(piNatives.Shell.prototype, "run").mockImplementation((options, onChunk) => { + if (options.signal instanceof AbortSignal) { + nativeSignal = options.signal; + } + const nativeResult = Promise.withResolvers(); + realSetTimeout(() => { + onChunk?.(null, "flushed-during-timeout\n"); + nativeResult.resolve({ exitCode: undefined, cancelled: false, timedOut: true }); + }, 20); + return nativeResult.promise; + }); + const abortSpy = vi.spyOn(piNatives.Shell.prototype, "abort").mockResolvedValue(); + + const result = await executeBash("producer | tail -5", { + cwd: tempDir, + timeout: 1000, + sessionKey: "native-timeout-flushes-pipeline", + }); + + expect(result.cancelled).toBe(true); + expect(result.output).toContain("flushed-during-timeout"); + expect(result.output).toContain("Command timed out after 1 seconds"); + expect(nativeSignal).toBeDefined(); + expect(nativeSignal?.aborted).toBe(false); + expect(abortSpy).not.toHaveBeenCalled(); + }); + + it("keeps a delayed JavaScript fallback for stalled native timeout cleanup", async () => { + const realSetTimeout = globalThis.setTimeout; + let fallbackDelayMs = 0; + vi.spyOn(globalThis, "setTimeout").mockImplementation(((handler: () => void, ms?: number, ...rest: unknown[]) => { + if (typeof ms === "number" && ms >= 1000) { + fallbackDelayMs = Math.max(fallbackDelayMs, ms); + return realSetTimeout(handler, 5, ...rest); + } + return realSetTimeout(handler, ms, ...rest); + }) as typeof globalThis.setTimeout); + let nativeSignal: AbortSignal | undefined; vi.spyOn(piNatives.Shell.prototype, "run").mockImplementation((options, onChunk) => { if (options.signal instanceof AbortSignal) { @@ -552,6 +587,7 @@ exit 64 expect(result.cancelled).toBe(true); expect(result.output).toContain("streamed-before-timeout"); expect(result.output).toContain("Command timed out after 1 seconds"); + expect(fallbackDelayMs).toBeGreaterThan(1000); expect(nativeSignal).toBeDefined(); expect(nativeSignal?.aborted).toBe(false); expect(abortSpy).not.toHaveBeenCalled(); diff --git a/packages/natives/CHANGELOG.md b/packages/natives/CHANGELOG.md index 7f715115d..b18a474de 100644 --- a/packages/natives/CHANGELOG.md +++ b/packages/natives/CHANGELOG.md @@ -2,6 +2,10 @@ ## [Unreleased] +### Fixed + +- Fixed timed-out shell pipelines cancelling their output reader while the final stage was still flushing, which dropped captured output and could terminate Windows hosts during teardown. ([#5316](https://github.com/can1357/oh-my-pi/issues/5316)) + ## [16.4.6] - 2026-07-12 ### Added From 8d9ba57c1af776615180a8a584b0095b00d917bf Mon Sep 17 00:00:00 2001 From: can1357 Date: Fri, 17 Jul 2026 04:09:37 +0200 Subject: [PATCH 2/2] test(bash): keep timeout regression additive --- .../coding-agent/test/bash-executor.test.ts | 84 +++++++++---------- 1 file changed, 42 insertions(+), 42 deletions(-) diff --git a/packages/coding-agent/test/bash-executor.test.ts b/packages/coding-agent/test/bash-executor.test.ts index 06ce9dbe1..a6eec37ca 100644 --- a/packages/coding-agent/test/bash-executor.test.ts +++ b/packages/coding-agent/test/bash-executor.test.ts @@ -520,54 +520,19 @@ exit 64 expect(next.output.trim()).toBe("still_persistent"); }); - it("waits for native timeout teardown to flush piped output", async () => { + it("does not abort the native signal when the JavaScript timeout fallback returns streamed output", async () => { + // Compress the JS-side fallback timer (floored at 1000ms in the source) so + // the safety-net fires deterministically without a real 1s wait. Only long + // timers are shrunk — fs/subprocess setup keeps real scheduling — and the + // reported "1 seconds" derives from the configured timeout, not the timer. const realSetTimeout = globalThis.setTimeout; vi.spyOn(globalThis, "setTimeout").mockImplementation(((handler: () => void, ms?: number, ...rest: unknown[]) => realSetTimeout( handler, - ms === 1000 ? 5 : typeof ms === "number" && ms > 1000 ? 50 : ms, + typeof ms === "number" && ms >= 1000 ? 5 : ms, ...rest, )) as typeof globalThis.setTimeout); - let nativeSignal: AbortSignal | undefined; - vi.spyOn(piNatives.Shell.prototype, "run").mockImplementation((options, onChunk) => { - if (options.signal instanceof AbortSignal) { - nativeSignal = options.signal; - } - const nativeResult = Promise.withResolvers(); - realSetTimeout(() => { - onChunk?.(null, "flushed-during-timeout\n"); - nativeResult.resolve({ exitCode: undefined, cancelled: false, timedOut: true }); - }, 20); - return nativeResult.promise; - }); - const abortSpy = vi.spyOn(piNatives.Shell.prototype, "abort").mockResolvedValue(); - - const result = await executeBash("producer | tail -5", { - cwd: tempDir, - timeout: 1000, - sessionKey: "native-timeout-flushes-pipeline", - }); - - expect(result.cancelled).toBe(true); - expect(result.output).toContain("flushed-during-timeout"); - expect(result.output).toContain("Command timed out after 1 seconds"); - expect(nativeSignal).toBeDefined(); - expect(nativeSignal?.aborted).toBe(false); - expect(abortSpy).not.toHaveBeenCalled(); - }); - - it("keeps a delayed JavaScript fallback for stalled native timeout cleanup", async () => { - const realSetTimeout = globalThis.setTimeout; - let fallbackDelayMs = 0; - vi.spyOn(globalThis, "setTimeout").mockImplementation(((handler: () => void, ms?: number, ...rest: unknown[]) => { - if (typeof ms === "number" && ms >= 1000) { - fallbackDelayMs = Math.max(fallbackDelayMs, ms); - return realSetTimeout(handler, 5, ...rest); - } - return realSetTimeout(handler, ms, ...rest); - }) as typeof globalThis.setTimeout); - let nativeSignal: AbortSignal | undefined; vi.spyOn(piNatives.Shell.prototype, "run").mockImplementation((options, onChunk) => { if (options.signal instanceof AbortSignal) { @@ -587,7 +552,6 @@ exit 64 expect(result.cancelled).toBe(true); expect(result.output).toContain("streamed-before-timeout"); expect(result.output).toContain("Command timed out after 1 seconds"); - expect(fallbackDelayMs).toBeGreaterThan(1000); expect(nativeSignal).toBeDefined(); expect(nativeSignal?.aborted).toBe(false); expect(abortSpy).not.toHaveBeenCalled(); @@ -1025,6 +989,42 @@ exit 64 expect(result.output).toContain("Command cancelled"); await expectMarkerNeverWritten(marker, release); }); + it("waits for native timeout teardown to flush piped output", async () => { + const realSetTimeout = globalThis.setTimeout; + vi.spyOn(globalThis, "setTimeout").mockImplementation(((handler: () => void, ms?: number, ...rest: unknown[]) => + realSetTimeout( + handler, + ms === 1000 ? 5 : typeof ms === "number" && ms > 1000 ? 50 : ms, + ...rest, + )) as typeof globalThis.setTimeout); + + let nativeSignal: AbortSignal | undefined; + vi.spyOn(piNatives.Shell.prototype, "run").mockImplementation((options, onChunk) => { + if (options.signal instanceof AbortSignal) { + nativeSignal = options.signal; + } + const nativeResult = Promise.withResolvers(); + realSetTimeout(() => { + onChunk?.(null, "flushed-during-timeout\n"); + nativeResult.resolve({ exitCode: undefined, cancelled: false, timedOut: true }); + }, 20); + return nativeResult.promise; + }); + const abortSpy = vi.spyOn(piNatives.Shell.prototype, "abort").mockResolvedValue(); + + const result = await executeBash("producer | tail -5", { + cwd: tempDir, + timeout: 1000, + sessionKey: "native-timeout-flushes-pipeline", + }); + + expect(result.cancelled).toBe(true); + expect(result.output).toContain("flushed-during-timeout"); + expect(result.output).toContain("Command timed out after 1 seconds"); + expect(nativeSignal).toBeDefined(); + expect(nativeSignal?.aborted).toBe(false); + expect(abortSpy).not.toHaveBeenCalled(); + }); }); describe("executeBash :async: background retention", () => {