diff --git a/packages/core/src/process.ts b/packages/core/src/process.ts index 8d13d79433..a9bfb85fc3 100644 --- a/packages/core/src/process.ts +++ b/packages/core/src/process.ts @@ -9,6 +9,14 @@ export class AppProcessError extends Schema.TaggedErrorClass()( command: Schema.String, exitCode: Schema.optional(Schema.Number), stderr: Schema.optional(Schema.String), + /** + * Whatever the child had written to stdout when this failed. + * + * Present so a TIMEOUT says something. Without it the only report was "Timed out", which cannot tell a + * child that hung before producing a byte from one that did most of its work and then stalled -- and on + * a flake that only reproduces on one platform in CI, that distinction is the whole investigation. + */ + stdout: Schema.optional(Schema.String), cause: Schema.optional(Schema.Defect()), }) { override get message() { @@ -118,10 +126,30 @@ const normalizeStdin = ( ? Stream.make(input) : input -export const collectStream = (stream: Stream.Stream, maxOutputBytes: number | undefined) => +/** + * What `collectStream` gathers. Exposed so a CALLER can own it. + * + * The fold mutates one object and hands the same one back on every chunk, which means a caller holding a + * reference sees whatever arrived even if the fold never completes. That is the whole point: a command + * killed by its timeout used to report nothing at all about the child, because the buffers lived inside + * the interrupted fiber and died with it. A subprocess that fails by producing NOTHING and one that fails + * after printing half its work are very different diagnoses, and both looked identical. + */ +export type StreamAccumulator = { chunks: Uint8Array[]; bytes: number; truncated: boolean } + +export const newAccumulator = (): StreamAccumulator => ({ chunks: [], bytes: 0, truncated: false }) + +/** Whatever the accumulator holds right now, as a buffer. Safe to call mid-stream. */ +export const accumulated = (acc: StreamAccumulator) => Buffer.concat(acc.chunks) + +export const collectStream = ( + stream: Stream.Stream, + maxOutputBytes: number | undefined, + into: StreamAccumulator = newAccumulator(), +) => Stream.runFold( stream, - () => ({ chunks: [] as Uint8Array[], bytes: 0, truncated: false }), + () => into, (acc, chunk) => { if (maxOutputBytes === undefined) { acc.chunks.push(chunk) @@ -143,12 +171,19 @@ const layer = Layer.effect( const runCommand = (command: ChildProcess.Command, options?: RunOptions) => { const description = describeCommand(command) + /* + Owned HERE, outside `collect`, so a timeout can still read them. `Effect.timeoutOrElse` interrupts + `collect`, and anything scoped to that fiber goes with it -- which is why a timed-out command used + to report neither stream. + */ + const outAcc = newAccumulator() + const errAcc = newAccumulator() const collect = Effect.scoped( Effect.gen(function* () { const handle = yield* spawner.spawn(command) if (options?.combineOutput) { const [output, exitCode] = yield* Effect.all( - [collectStream(handle.all, options.maxOutputBytes), handle.exitCode], + [collectStream(handle.all, options.maxOutputBytes, outAcc), handle.exitCode], { concurrency: "unbounded" }, ) return { @@ -164,8 +199,8 @@ const layer = Layer.effect( } const [stdout, stderr, exitCode] = yield* Effect.all( [ - collectStream(handle.stdout, options?.maxOutputBytes), - collectStream(handle.stderr, options?.maxErrorBytes), + collectStream(handle.stdout, options?.maxOutputBytes, outAcc), + collectStream(handle.stderr, options?.maxErrorBytes, errAcc), handle.exitCode, ], { concurrency: "unbounded" }, @@ -183,7 +218,19 @@ const layer = Layer.effect( const timed = options?.timeout ? Effect.timeoutOrElse(collect, { duration: options.timeout, - orElse: () => Effect.fail(new AppProcessError({ command: description, cause: new Error("Timed out") })), + orElse: () => + Effect.fail( + new AppProcessError({ + command: description, + /* + Partial output, not nothing. An empty stdout here is now a real observation about the + child rather than an artifact of how the failure was built. + */ + stdout: accumulated(outAcc).toString(), + stderr: accumulated(errAcc).toString(), + cause: new Error("Timed out"), + }), + ), }) : collect const aborted = options?.signal diff --git a/packages/core/test/process/process.test.ts b/packages/core/test/process/process.test.ts index b213ab1559..afa278f49a 100644 --- a/packages/core/test/process/process.test.ts +++ b/packages/core/test/process/process.test.ts @@ -3,7 +3,7 @@ import fs from "fs/promises" import { realpathSync } from "node:fs" import { tmpdir } from "node:os" import path from "node:path" -import { Effect, Exit, Fiber, Stream } from "effect" +import { Cause, Effect, Exit, Fiber, Stream } from "effect" import { ChildProcess } from "effect/unstable/process" import { LayerNode } from "@redrob-code/core/effect/layer-node" import { AppProcess } from "@redrob-code/core/process" @@ -152,6 +152,31 @@ describe("AppProcess", () => { ) if (process.platform !== "win32") { + /* + The claim this pins: a command killed by its timeout reports what the child had ALREADY printed. + + It used to report nothing on either stream, because the accumulators lived inside the fiber + `Effect.timeoutOrElse` interrupts. That turned every subprocess timeout into the same message + regardless of cause -- a child that hung before writing a byte and one that did most of its work + and then stalled were indistinguishable, which is exactly the distinction an intermittent + platform-specific hang turns on. + */ + it.live( + "a timeout reports the output the child produced before it was killed", + Effect.gen(function* () { + const svc = yield* AppProcess.Service + /* Prints, flushes, then hangs -- so there IS partial output to lose. */ + const script = `process.stdout.write('half-done\\n');process.stderr.write('warned\\n');setInterval(()=>{},60000)` + const exit = yield* Effect.exit(svc.run(cmd("-e", script), { timeout: "700 millis" })) + expect(Exit.isFailure(exit)).toBe(true) + if (!Exit.isFailure(exit)) return + const error = Cause.squash(exit.cause) as { stdout?: string; stderr?: string } + expect(error.stdout).toContain("half-done") + expect(error.stderr).toContain("warned") + }), + 8_000, + ) + it.live( "timeout cleans up the scoped child process", Effect.acquireUseRelease( diff --git a/packages/redrob/test/cli/acp/acp-test-client.ts b/packages/redrob/test/cli/acp/acp-test-client.ts index 1588d344d7..f8bd4ae2af 100644 --- a/packages/redrob/test/cli/acp/acp-test-client.ts +++ b/packages/redrob/test/cli/acp/acp-test-client.ts @@ -2,6 +2,7 @@ import { expect } from "bun:test" import type { SessionConfigOption, SessionConfigSelectOption } from "@agentclientprotocol/sdk" import { Duration, Effect } from "effect" import type { AcpHandle } from "../../lib/cli-process" +import { slowPlatform } from "../../lib/cli-process" type JsonRpcRequest = { readonly jsonrpc: "2.0" @@ -44,15 +45,19 @@ export function createAcpClient(acp: AcpHandle): AcpClient { yield* acp.send(message) while (true) { - const received = yield* acp.receive.pipe(Effect.timeout(Duration.seconds(15))) + const received = yield* acp.receive.pipe(Effect.timeout(Duration.millis(slowPlatform(15_000)))) if (isJsonRpcResponse(received) && received.id === id) return received } }) + /* + The default is scaled for the slow platform, as are the explicit values callers pass -- every wait in + this client is on a subprocess round trip, and Windows is where these bounds flake. + */ const waitForNotification = (method: string, predicate: (params: T) => boolean, timeoutMs = 15_000) => Effect.gen(function* () { while (true) { - const received = yield* acp.receive.pipe(Effect.timeout(Duration.millis(timeoutMs))) + const received = yield* acp.receive.pipe(Effect.timeout(Duration.millis(slowPlatform(timeoutMs)))) if (!isJsonRpcNotification(received)) continue if (received.method === method && predicate(received.params as T)) return received } diff --git a/packages/redrob/test/cli/acp/lifecycle.test.ts b/packages/redrob/test/cli/acp/lifecycle.test.ts index eac2693f1e..de606bb0f0 100644 --- a/packages/redrob/test/cli/acp/lifecycle.test.ts +++ b/packages/redrob/test/cli/acp/lifecycle.test.ts @@ -6,7 +6,7 @@ import type { ResumeSessionResponse, } from "@agentclientprotocol/sdk" import { Duration, Effect } from "effect" -import { cliIt } from "../../lib/cli-process" +import { cliIt, slowPlatform } from "../../lib/cli-process" import { expectOk, selectConfigOption } from "./acp-test-client" import { createAcpClient, initialize, newSession, verifierConfig } from "./helpers" @@ -18,7 +18,10 @@ describe("redrob acp lifecycle subprocess", () => { const acp = yield* redrob.acp() acp.close() - const code = yield* Effect.promise(() => acp.exited).pipe(Effect.timeout(Duration.seconds(5))) + /* Scaled: Windows process teardown is slower, and five seconds was tight enough to flake. */ + const code = yield* Effect.promise(() => acp.exited).pipe( + Effect.timeout(Duration.millis(slowPlatform(5_000))), + ) expect(code).toBe(0) }), 60_000, diff --git a/packages/redrob/test/cli/run/run-process.test.ts b/packages/redrob/test/cli/run/run-process.test.ts index cb318b5fea..f72728da03 100644 --- a/packages/redrob/test/cli/run/run-process.test.ts +++ b/packages/redrob/test/cli/run/run-process.test.ts @@ -6,7 +6,7 @@ import { describe, expect } from "bun:test" import { Effect } from "effect" import { reply } from "../../lib/llm-server" -import { cliIt } from "../../lib/cli-process" +import { cliIt, slowPlatform } from "../../lib/cli-process" describe("redrob run (non-interactive subprocess)", () => { // Happy path: prompt completes, output reaches stdout, process exits 0. @@ -73,10 +73,12 @@ describe("redrob run (non-interactive subprocess)", () => { Effect.gen(function* () { const result = yield* redrob.run("say hi", { model: "test/nonexistent-model", - timeoutMs: 15_000, + timeoutMs: slowPlatform(15_000), }) expect(result.exitCode).not.toBe(0) - expect(result.durationMs).toBeLessThan(15_000) + /* Scaled WITH the bound: the claim is that it exits promptly rather than being killed, and a + duration assertion pinned to a base figure while the bound moves tests the wrong thing. */ + expect(result.durationMs).toBeLessThan(slowPlatform(15_000)) }), 60_000, ) @@ -96,7 +98,7 @@ describe("redrob run (non-interactive subprocess)", () => { ) yield* llm.fail("upstream provider exploded mid-stream") yield* llm.text("recovered") - const result = yield* redrob.run("trigger midstream error", { timeoutMs: 30_000 }) + const result = yield* redrob.run("trigger midstream error", { timeoutMs: slowPlatform(30_000) }) expect(result.exitCode).toBe(0) expect(result.stdout).toBe("partial response\nrecovered\n") expect(result.stderr).not.toContain("upstream provider exploded mid-stream") diff --git a/packages/redrob/test/lib/cli-process.ts b/packages/redrob/test/lib/cli-process.ts index eee5542c3c..80537c6d0e 100644 --- a/packages/redrob/test/lib/cli-process.ts +++ b/packages/redrob/test/lib/cli-process.ts @@ -42,6 +42,37 @@ const cliEntry = path.join(redrobRoot, "src/index.ts") // whole test and would hold permits while a nested `run` waits for one. const spawnGate = Semaphore.makeUnsafe(Math.max(2, availableParallelism() - 1)) +/** + * The harness's wait bounds, derived from one number instead of written three times. + * + * They are not independent: `cliIt.concurrent` carried a comment saying Bun's test timeout must stay + * ABOVE the child timeout, while the two numbers sat in different functions as unrelated literals with + * nothing keeping them in step. Raising one without the other silently breaks the invariant the comment + * promises, and the failure mode is a test that expires holding no spawn permit -- which reads as the + * command being slow rather than as the harness mis-tuned. + * + * Windows gets a multiple of the base. A `redrob.spawn` is `bun run` over the TypeScript entry, so every + * spawn pays a cold transpile plus that platform's process-creation cost, and the runners are slower + * again. This is a MITIGATION and not a diagnosis: the flake it responds to timed out at exactly the + * bound with no output preserved, which is now fixed separately -- the next occurrence will say where the + * child actually got to, and this number should be revisited against that evidence rather than raised + * again by feel. + */ +const SLOW_PLATFORM_FACTOR = process.platform === "win32" ? 3 : 1 +const CHILD_TIMEOUT_MS = 30_000 * SLOW_PLATFORM_FACTOR +const READY_TIMEOUT_MS = 15_000 * SLOW_PLATFORM_FACTOR +const TEST_TIMEOUT_MS = CHILD_TIMEOUT_MS * 2 + +/** + * A test's own bound, scaled for the slow platform. + * + * For the call sites that pass an explicit `timeoutMs` rather than taking the default -- those bypass the + * allowance entirely, which is how the first version of this fix left the very test that was flaking still + * pinned to thirty seconds. Write `slowPlatform(30_000)` and the intent stays readable while the platform + * correction is applied for you. + */ +export const slowPlatform = (ms: number) => ms * SLOW_PLATFORM_FACTOR + export const testModelID = "test/test-model" // Wrap a Bun subprocess pipe (or any ReadableStream) as a Stream. @@ -219,7 +250,7 @@ export function withCliFixture( return yield* spawnGate.withPermit( Effect.gen(function* () { const start = Date.now() - const timeoutMs = opts?.timeoutMs ?? 30_000 + const timeoutMs = opts?.timeoutMs ?? CHILD_TIMEOUT_MS // stdin: "ignore" so the child doesn't see a piped stdin and block // on `Bun.stdin.text()` (see src/cli/cmd/run.ts — non-TTY stdin is // consumed as the prompt). The old Process.run wrapper defaulted to @@ -246,7 +277,13 @@ export function withCliFixture( Effect.succeed({ command: err.command, exitCode: err.exitCode ?? -1, - stdout: Buffer.alloc(0), + /* + The child's own output, not an empty buffer. This used to be `Buffer.alloc(0)` + unconditionally, so `expectExit`'s dump printed an empty stdout on every timeout and a + reader could not tell whether the child had produced nothing or the harness had thrown it + away. It was the latter, which cost a real investigation. + */ + stdout: Buffer.from(err.stdout ?? ""), stderr: Buffer.from((err.stderr ?? String(err.cause ?? err.message)) + "\n"), stdoutTruncated: false, stderrTruncated: false, @@ -375,7 +412,7 @@ export function withCliFixture( ), ) - const readyTimeoutMs = opts?.readyTimeoutMs ?? 15_000 + const readyTimeoutMs = opts?.readyTimeoutMs ?? READY_TIMEOUT_MS const match = yield* Deferred.await(readyDeferred).pipe( Effect.timeoutOrElse({ duration: Duration.millis(readyTimeoutMs), @@ -545,8 +582,16 @@ export const cliIt = { (process.platform === "win32" ? test : test.concurrent)( name, () => Effect.runPromise(Effect.scoped(withCliFixture(body))), - // Bun's timeout includes spawn-gate wait. Keep it above the child timeout so - // queued concurrent CLI tests do not expire while holding no permit. - opts ?? 60_000, + /* + A caller's own number is a FLOOR, not a ceiling. + + Bun's timeout includes spawn-gate wait, so it must stay above the child timeout or a queued test + expires while holding no permit -- and that failure reads as a slow command rather than as a + mis-tuned harness. Dozens of call sites pass `60_000`, which was comfortably above the old fixed + 30s child bound and is BELOW the Windows one, so honouring them literally would have reintroduced + exactly the drift these constants exist to prevent. Clamped here rather than edited at every call + site: the invariant then holds by construction instead of by everyone remembering it. + */ + typeof opts === "number" ? Math.max(opts, TEST_TIMEOUT_MS) : (opts ?? TEST_TIMEOUT_MS), ), }