diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml new file mode 100644 index 0000000..212bb77 --- /dev/null +++ b/.github/workflows/ci.yml @@ -0,0 +1,33 @@ +name: CI + +on: + pull_request: + branches: + - beta + push: + branches: + - beta + +permissions: + contents: read + +jobs: + verify: + runs-on: ubuntu-latest + timeout-minutes: 20 + steps: + - uses: actions/checkout@v4 + + - uses: oven-sh/setup-bun@v2 + with: + bun-version: 1.3.14 + + - uses: actions/setup-node@v4 + with: + node-version: 22 + + - name: Install exact dependencies + run: bun install --frozen-lockfile + + - name: Verify V2 beta package + run: npm run verify diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml index eea4fca..3da01f5 100644 --- a/.github/workflows/release.yml +++ b/.github/workflows/release.yml @@ -1,136 +1,90 @@ -name: Release +name: Release V2 Beta on: workflow_dispatch: inputs: - version: - description: 'Explicit version (e.g., 0.3.0). If empty, bump type is used.' - required: false - type: string - bump: - description: 'Version bump type (used if version is empty)' - required: false - default: 'minor' - type: choice - options: - - minor - - patch - - major dry_run: - description: 'Dry run — skip npm publish, create draft release' - required: false - default: false + description: Verify and pack without publishing + required: true + default: true type: boolean -env: - NODE_VERSION: '22' +concurrency: + group: opencode-cursor-v2-release + cancel-in-progress: false permissions: contents: write + id-token: write jobs: release: + if: github.ref_name == 'beta' runs-on: ubuntu-latest + timeout-minutes: 25 steps: - uses: actions/checkout@v4 + + - uses: oven-sh/setup-bun@v2 with: - fetch-depth: 0 - fetch-tags: true + bun-version: 1.3.14 - - name: Get current version from package.json - id: current_version - run: echo "version=$(node -p "require('./package.json').version")" >> "$GITHUB_OUTPUT" + - uses: actions/setup-node@v4 + with: + node-version: 22 + registry-url: https://registry.npmjs.org - - name: Calculate new version + - name: Validate prerelease version id: version - env: - CURRENT_VERSION: ${{ steps.current_version.outputs.version }} - INPUT_VERSION: ${{ github.event.inputs.version }} - BUMP_TYPE: ${{ github.event.inputs.bump }} + shell: bash run: | set -euo pipefail - if [[ -n "$INPUT_VERSION" ]]; then - # Validate semver - if ! echo "$INPUT_VERSION" | grep -qE '^[0-9]+\.[0-9]+\.[0-9]+$'; then - echo "Error: version must be in semver format X.Y.Z" - exit 1 - fi - echo "version=$INPUT_VERSION" >> "$GITHUB_OUTPUT" - else - IFS='.' read -r major minor patch <<< "$CURRENT_VERSION" - case "$BUMP_TYPE" in - major) echo "version=$((major+1)).0.0" >> "$GITHUB_OUTPUT" ;; - minor) echo "version=$major.$((minor+1)).0" >> "$GITHUB_OUTPUT" ;; - patch) echo "version=$major.$minor.$((patch+1))" >> "$GITHUB_OUTPUT" ;; - *) echo "Error: unknown bump type $BUMP_TYPE"; exit 1 ;; - esac + version="$(node -p "require('./package.json').version")" + if [[ ! "$version" =~ ^[0-9]+\.[0-9]+\.[0-9]+-beta\.[0-9]+$ ]]; then + echo "Expected package version X.Y.Z-beta.N, got: $version" >&2 + exit 1 fi + echo "version=$version" >> "$GITHUB_OUTPUT" - - name: Show version info - run: | - echo "Current version: ${{ steps.current_version.outputs.version }}" - echo "New version: ${{ steps.version.outputs.version }}" - echo "Dry run: ${{ github.event.inputs.dry_run == 'true' && 'yes' || 'no' }}" - - - name: Update package.json version - run: | - node -e " - const pkg = require('./package.json'); - pkg.version = '${{ steps.version.outputs.version }}'; - require('fs').writeFileSync('./package.json', JSON.stringify(pkg, null, 2) + '\n'); - " - - - name: Setup Bun - uses: oven-sh/setup-bun@v2 - - - name: Setup Node.js - uses: actions/setup-node@v4 - with: - node-version: ${{ env.NODE_VERSION }} - registry-url: 'https://registry.npmjs.org' - - - name: Install dependencies + - name: Install exact dependencies run: bun install --frozen-lockfile - - name: Run tests - run: bun run test + - name: Verify source and installed artifact + run: npm run verify - - name: Build package - run: bun run build - - - name: Commit and tag - env: - VERSION: ${{ steps.version.outputs.version }} + - name: Pack exact release artifact + id: pack + shell: bash run: | - git config user.name "github-actions[bot]" - git config user.email "github-actions[bot]@users.noreply.github.com" - git add package.json - if git ls-files --modified --others | grep -q bun.lock; then - git add bun.lock - fi - git commit -m "chore: release v${VERSION}" - git tag "v${VERSION}" - git push origin "v${VERSION}" - git push origin HEAD:main + set -euo pipefail + output="$(npm pack --ignore-scripts --json --pack-destination "$RUNNER_TEMP")" + filename="$(node -e 'const fs=require("fs"); const x=JSON.parse(fs.readFileSync(0,"utf8")); process.stdout.write(x[0].filename)' <<<"$output")" + echo "path=$RUNNER_TEMP/$filename" >> "$GITHUB_OUTPUT" - - name: Publish to npm - if: ${{ github.event.inputs.dry_run != 'true' }} - run: npm publish --access public + - name: Upload verified artifact + uses: actions/upload-artifact@v4 + with: + name: opencode-cursor-v2-${{ steps.version.outputs.version }} + path: ${{ steps.pack.outputs.path }} + + - name: Publish verified artifact + if: inputs.dry_run == false + run: npm publish "${{ steps.pack.outputs.path }}" --tag beta --provenance --access public env: NODE_AUTH_TOKEN: ${{ secrets.NPM_TOKEN }} - - name: Create GitHub Release + - name: Create GitHub prerelease + if: inputs.dry_run == false uses: softprops/action-gh-release@v2 with: tag_name: v${{ steps.version.outputs.version }} + target_commitish: ${{ github.sha }} generate_release_notes: true - draft: ${{ github.event.inputs.dry_run == 'true' }} + prerelease: true name: v${{ steps.version.outputs.version }} - env: - GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }} - name: Send Discord notification - if: ${{ github.event.inputs.dry_run != 'true' }} + if: inputs.dry_run == false env: DISCORD_WEBHOOK_URL: ${{ secrets.DISCORD_WEBHOOK_URL }} GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }} diff --git a/package.json b/package.json index f23dcd6..475c96f 100644 --- a/package.json +++ b/package.json @@ -25,7 +25,7 @@ "scripts": { "build": "node scripts/clean.mjs && tsc -p tsconfig.json && node scripts/copy-runtime.mjs", "test": "bun test/smoke.ts", - "test:v2": "bun test test/v2-*.test.ts test/node-runtime.test.ts test/model-normalizer.test.ts", + "test:v2": "bun test test/v2-*.test.ts test/bridge-pool.test.ts test/node-runtime.test.ts test/model-normalizer.test.ts test/shared-constants.test.ts", "test:package": "node scripts/check-package.mjs", "test:v2-loader": "node scripts/smoke-opencode-v2.mjs", "typecheck": "tsc -p tsconfig.json --noEmit", diff --git a/src/bridge-pool.ts b/src/bridge-pool.ts index eb4736e..0fe3240 100644 --- a/src/bridge-pool.ts +++ b/src/bridge-pool.ts @@ -11,6 +11,7 @@ */ import { fileURLToPath } from "node:url"; import { resolveNodeExecutable } from "./node-runtime.js"; +import { log } from "./shared/log.js"; const PERSISTENT_BRIDGE_PATH = fileURLToPath( new URL("./h2-bridge-persistent.mjs", import.meta.url), @@ -176,6 +177,13 @@ export interface BridgeAcquireOptions { unary?: boolean; } +export class BridgePoolCapacityError extends Error { + constructor(maxSize: number) { + super(`BridgePool capacity reached (${maxSize} active workers)`); + this.name = "BridgePoolCapacityError"; + } +} + export class BridgePool { private idle: PersistentWorker[] = []; private allWorkers = new Set(); @@ -184,15 +192,21 @@ export class BridgePool { private shuttingDown = false; constructor(options: BridgePoolOptions = {}) { - this.minSize = options.minSize ?? 2; - this.maxSize = options.maxSize ?? 4; + const minSize = options.minSize ?? 2; + const maxSize = options.maxSize ?? 4; + if (!Number.isSafeInteger(minSize) || minSize < 0) { + throw new Error(`BridgePool minSize must be a non-negative integer, got ${minSize}`); + } + if (!Number.isSafeInteger(maxSize) || maxSize < 1) { + throw new Error(`BridgePool maxSize must be a positive integer, got ${maxSize}`); + } + this.minSize = Math.min(minSize, maxSize); + this.maxSize = maxSize; } /** Pre-warm the pool with minSize idle workers. */ warmup(): void { - for (let i = 0; i < this.minSize; i++) { - this.addWorker(); - } + this.replenish(); } /** Acquire a bridge handle for a new request. */ @@ -217,9 +231,7 @@ export class BridgePool { worker = this.addWorker(); this.idle.pop(); } else { - // Pool full — spawn an ephemeral worker not tracked by pool - worker = spawnWorker(); - // Don't add to allWorkers — it won't be returned to pool + throw new BridgePoolCapacityError(this.maxSize); } } @@ -305,7 +317,6 @@ export class BridgePool { /** Exit code recorded when STREAM_DONE/process-death completes before callers attach onClose. */ let recordedExitCode = 0; const pool = this; - const isPooled = this.allWorkers.has(worker); /** Buffer OUTPUT_DATA until the caller registers onData (stdout can beat ReadableStream wiring). */ const pendingData: Buffer[] = []; let userDataCb: ((chunk: Buffer) => void) | null = null; @@ -318,6 +329,16 @@ export class BridgePool { // When stream completes (bridge sends STREAM_DONE), fire onClose and return to pool let closeCb: ((code: number) => void) | null = null; + const notifyClose = (callback: ((code: number) => void) | null, code: number) => { + if (!callback) return; + queueMicrotask(() => { + try { + callback(code); + } catch (error) { + log.error("[bridge-pool] onClose callback failed", error); + } + }); + }; worker.cbs.streamDone = (code: number) => { if (done) return; @@ -325,16 +346,8 @@ export class BridgePool { recordedExitCode = code; const cbNow = closeCb; closeCb = null; - cbNow?.(code); - if (isPooled) { - pool.release(worker); - } else { - // Ephemeral overflow worker (pool was saturated at acquire time): it is - // not tracked by the pool and will never be reused, so shut it down - // instead of leaking the child process. - workerSendShutdown(worker); - workerKill(worker); - } + pool.release(worker); + notifyClose(cbNow, code); }; // Handle unexpected process death @@ -344,12 +357,8 @@ export class BridgePool { recordedExitCode = 1; const cbNow = closeCb; closeCb = null; - cbNow?.(1); - if (isPooled) { - pool.remove(worker); - } else { - workerKill(worker); - } + pool.remove(worker); + notifyClose(cbNow, 1); }; return { @@ -366,11 +375,11 @@ export class BridgePool { kill() { if (done) return; done = true; - if (isPooled) { - pool.remove(worker); - } else { - workerKill(worker); - } + recordedExitCode = 1; + const cbNow = closeCb; + closeCb = null; + pool.remove(worker); + notifyClose(cbNow, 1); }, onData(cb: (chunk: Buffer) => void) { const flushed = pendingData.splice(0, pendingData.length); @@ -379,7 +388,7 @@ export class BridgePool { }, onClose(cb: (code: number) => void) { if (done) { - queueMicrotask(() => cb(recordedExitCode)); + notifyClose(cb, recordedExitCode); } else { closeCb = cb; } diff --git a/src/proxy.ts b/src/proxy.ts index 12f39d9..bb1b225 100644 --- a/src/proxy.ts +++ b/src/proxy.ts @@ -112,7 +112,11 @@ import { deterministicConversationId, selectionIdentity, } from "./conversation/identity.js"; -import { BridgePool, type BridgeHandle } from "./bridge-pool.js"; +import { + BridgePool, + BridgePoolCapacityError, + type BridgeHandle, +} from "./bridge-pool.js"; import { log } from "./shared/log.js"; import { CURSOR_SELECTION_HEADER, @@ -801,6 +805,24 @@ export async function startProxy( if (isAbortError(err) || req.signal.aborted) { return new Response(null, { status: 499, statusText: "Client Closed Request" }); } + if (err instanceof BridgePoolCapacityError) { + return new Response( + JSON.stringify({ + error: { + message: "Server is saturated, please retry shortly", + type: "server_error", + code: "service_unavailable", + }, + }), + { + status: 503, + headers: { + "Content-Type": "application/json", + "Retry-After": "2", + }, + }, + ); + } const message = err instanceof Error ? err.message : String(err); return new Response( JSON.stringify({ @@ -3044,4 +3066,4 @@ function handleToolResultResume( false, abortSignal, ); -} \ No newline at end of file +} diff --git a/src/shared/constants.ts b/src/shared/constants.ts index 3072cb0..5d373aa 100644 --- a/src/shared/constants.ts +++ b/src/shared/constants.ts @@ -1,12 +1,23 @@ /** Shared plugin constants — single source of truth. */ +function positiveInteger(value: string | undefined, fallback: number): number { + const parsed = Math.floor(Number(value)); + return Number.isSafeInteger(parsed) && parsed > 0 ? parsed : fallback; +} + export const CURSOR_PROVIDER_ID = "cursor"; export const DEFAULT_MODEL_ID = "default"; export const OPENAI_COMPATIBLE_NPM = "@ai-sdk/openai-compatible"; export const CURSOR_VARIANT_OPTION = "cursorVariant"; -export const DEFAULT_CONTEXT_WINDOW = 200_000; -export const DEFAULT_MAX_TOKENS = 64_000; +export const DEFAULT_CONTEXT_WINDOW = positiveInteger( + process.env.OPENCODE_CURSOR_DEFAULT_CONTEXT_WINDOW, + 200_000, +); +export const DEFAULT_MAX_TOKENS = positiveInteger( + process.env.OPENCODE_CURSOR_DEFAULT_MAX_TOKENS, + 64_000, +); export const GENERATED_VARIANT_KEYS = [ "none", diff --git a/test/bridge-pool.test.ts b/test/bridge-pool.test.ts new file mode 100644 index 0000000..427b45c --- /dev/null +++ b/test/bridge-pool.test.ts @@ -0,0 +1,202 @@ +import { describe, expect, test } from "bun:test"; +import http2 from "node:http2"; +import type { AddressInfo } from "node:net"; +import { BridgePool, type BridgeHandle } from "../src/bridge-pool"; + +async function createServer() { + let streamCount = 0; + const sessions = new Set(); + const held = new Set(); + const server = http2.createServer(); + server.on("session", (session) => { + sessions.add(session); + session.once("close", () => sessions.delete(session)); + }); + server.on("stream", (stream, headers) => { + streamCount += 1; + stream.respond({ ":status": 200, "content-type": "application/connect+proto" }); + if (headers[":path"] === "/hold") { + held.add(stream); + stream.once("close", () => held.delete(stream)); + return; + } + stream.end(Buffer.from(`response-${streamCount}`)); + }); + await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); + return { + url: `http://127.0.0.1:${(server.address() as AddressInfo).port}`, + streamCount: () => streamCount, + releaseHeld() { + for (const stream of held) stream.end(); + }, + async close() { + for (const stream of held) stream.destroy(); + for (const session of sessions) session.destroy(); + await new Promise((resolve, reject) => { + server.close((error) => error ? reject(error) : resolve()); + }); + }, + }; +} + +function complete(handle: BridgeHandle): Promise<{ code: number; data: string }> { + return new Promise((resolve) => { + const chunks: Buffer[] = []; + handle.onData((chunk) => chunks.push(chunk)); + handle.onClose((code) => { + resolve({ code, data: Buffer.concat(chunks).toString("utf8") }); + }); + handle.end(); + }); +} + +describe("BridgePool", () => { + test("reuses a pooled worker across sequential requests", async () => { + const server = await createServer(); + const pool = new BridgePool({ minSize: 0, maxSize: 1 }); + try { + const first = await complete(pool.acquire({ + accessToken: "token", + rpcPath: "/run", + url: server.url, + })); + const second = await complete(pool.acquire({ + accessToken: "token", + rpcPath: "/run", + url: server.url, + })); + + expect(first.code).toBe(0); + expect(second.code).toBe(0); + expect(server.streamCount()).toBe(2); + expect(pool.stats()).toEqual({ idle: 1, active: 0, total: 1, maxSize: 1 }); + } finally { + pool.shutdown(); + await server.close(); + } + }); + + test("rejects new work instead of spawning an overflow worker", async () => { + const server = await createServer(); + const pool = new BridgePool({ minSize: 0, maxSize: 1 }); + try { + const active = pool.acquire({ + accessToken: "token", + rpcPath: "/hold", + url: server.url, + }); + active.onData(() => {}); + active.onClose(() => {}); + active.end(); + + expect(() => pool.acquire({ + accessToken: "token", + rpcPath: "/run", + url: server.url, + })).toThrow("capacity reached"); + expect(pool.stats()).toEqual({ idle: 0, active: 1, total: 1, maxSize: 1 }); + } finally { + server.releaseHeld(); + pool.shutdown(); + await server.close(); + } + }); + + test("notifies an active reader when its handle is killed", async () => { + const server = await createServer(); + const pool = new BridgePool({ minSize: 0, maxSize: 1 }); + try { + const active = pool.acquire({ + accessToken: "token", + rpcPath: "/hold", + url: server.url, + }); + active.onData(() => {}); + const closed = new Promise((resolve) => active.onClose(resolve)); + + active.kill(); + + expect(await closed).toBe(1); + expect(active.alive).toBe(false); + } finally { + pool.shutdown(); + await server.close(); + } + }); + + test("releases capacity before notifying a completed reader", async () => { + const server = await createServer(); + const pool = new BridgePool({ minSize: 0, maxSize: 1 }); + try { + const first = pool.acquire({ + accessToken: "token", + rpcPath: "/run", + url: server.url, + }); + first.onData(() => {}); + const second = new Promise<{ code: number; data: string }>((resolve, reject) => { + first.onClose(() => { + try { + complete(pool.acquire({ + accessToken: "token", + rpcPath: "/run", + url: server.url, + })).then(resolve, reject); + } catch (error) { + reject(error); + } + }); + }); + first.end(); + + expect((await second).code).toBe(0); + } finally { + pool.shutdown(); + await server.close(); + } + }); + + test("releases capacity when a close callback throws", async () => { + const server = await createServer(); + const pool = new BridgePool({ minSize: 0, maxSize: 1 }); + try { + const first = pool.acquire({ + accessToken: "token", + rpcPath: "/run", + url: server.url, + }); + first.onData(() => {}); + const notified = new Promise((resolve) => { + first.onClose(() => { + resolve(); + throw new Error("test callback failure"); + }); + }); + first.end(); + await notified; + + expect((await complete(pool.acquire({ + accessToken: "token", + rpcPath: "/run", + url: server.url, + }))).code).toBe(0); + } finally { + pool.shutdown(); + await server.close(); + } + }); + + test("keeps warmup within validated pool bounds", () => { + const pool = new BridgePool({ minSize: 2, maxSize: 1 }); + try { + pool.warmup(); + pool.warmup(); + expect(pool.stats()).toEqual({ idle: 1, active: 0, total: 1, maxSize: 1 }); + } finally { + pool.shutdown(); + } + + expect(() => new BridgePool({ minSize: -1 })).toThrow("minSize"); + expect(() => new BridgePool({ maxSize: Number.NaN })).toThrow("maxSize"); + }); +}); diff --git a/test/shared-constants.test.ts b/test/shared-constants.test.ts new file mode 100644 index 0000000..3978172 --- /dev/null +++ b/test/shared-constants.test.ts @@ -0,0 +1,39 @@ +import { describe, expect, test } from "bun:test"; +import { execFileSync } from "node:child_process"; + +const readFallbackLimits = (environment: NodeJS.ProcessEnv) => JSON.parse( + execFileSync( + process.execPath, + [ + "-e", + "import('./src/shared/constants.ts').then(({ DEFAULT_CONTEXT_WINDOW, DEFAULT_MAX_TOKENS }) => console.log(JSON.stringify([DEFAULT_CONTEXT_WINDOW, DEFAULT_MAX_TOKENS])))", + ], + { + cwd: process.cwd(), + encoding: "utf8", + env: environment, + }, + ), +) as [number, number]; + +describe("shared constants", () => { + test("uses positive integer fallback-limit overrides", () => { + expect(readFallbackLimits({ + ...process.env, + OPENCODE_CURSOR_DEFAULT_CONTEXT_WINDOW: "123456", + OPENCODE_CURSOR_DEFAULT_MAX_TOKENS: "7890", + })).toEqual([123456, 7890]); + + expect(readFallbackLimits({ + ...process.env, + OPENCODE_CURSOR_DEFAULT_CONTEXT_WINDOW: "invalid", + OPENCODE_CURSOR_DEFAULT_MAX_TOKENS: "-1", + })).toEqual([200000, 64000]); + + expect(readFallbackLimits({ + ...process.env, + OPENCODE_CURSOR_DEFAULT_CONTEXT_WINDOW: "0.5", + OPENCODE_CURSOR_DEFAULT_MAX_TOKENS: "7890.9", + })).toEqual([200000, 7890]); + }); +}); diff --git a/test/smoke.ts b/test/smoke.ts index 04f18cd..19e20e6 100644 --- a/test/smoke.ts +++ b/test/smoke.ts @@ -1793,13 +1793,9 @@ async function testPoolSequentialRequests() { console.log("[test] Pool sequential requests OK"); } -/** - * Test that concurrent requests beyond maxSize succeed via ephemeral overflow - * workers, and that those untracked workers do not inflate the tracked pool - * (regression guard for the ephemeral-worker leak fix in BridgePool). - */ -async function testPoolOverflowEphemeralWorkers() { - console.log("[test] Testing pool overflow ephemeral workers..."); +/** Test that maxSize bounds worker processes and released capacity is reusable. */ +async function testPoolCapacityBound() { + console.log("[test] Testing pool capacity bound..."); const { BridgePool } = await import("../src/bridge-pool"); const server = await createPoolTestServer(); @@ -1808,33 +1804,92 @@ async function testPoolOverflowEphemeralWorkers() { pool.warmup(); await new Promise((r) => setTimeout(r, 200)); - // Fire more concurrent requests than maxSize so the pool must spawn - // ephemeral overflow workers (not tracked in allWorkers). - const CONCURRENCY = 4; - const results = await Promise.all( - Array.from({ length: CONCURRENCY }, () => poolRequest(pool, server.url)), - ); - for (let i = 0; i < results.length; i++) { - assertEqual(results[i]!.code, 0, `Overflow request ${i} should succeed (code=0)`); + const active = pool.acquire({ + accessToken: "test-token", + rpcPath: "/agent.v1.AgentService/Run", + url: server.url, + }); + active.onData(() => {}); + const activeClosed = new Promise((resolve) => active.onClose(resolve)); + active.end(); + + let capacityError: unknown; + try { + pool.acquire({ + accessToken: "test-token", + rpcPath: "/agent.v1.AgentService/Run", + url: server.url, + }); + } catch (error) { + capacityError = error; } + assert( + capacityError instanceof Error && capacityError.message.includes("capacity reached"), + "Concurrent request beyond maxSize should fail fast", + ); - // Give streamDone-driven shutdown of ephemeral workers time to run. - await new Promise((r) => setTimeout(r, 200)); + assertEqual(await activeClosed, 0, "Active request should complete successfully"); + const afterRelease = await poolRequest(pool, server.url); + assertEqual(afterRelease.code, 0, "Released pool capacity should be reusable"); const stats = pool.stats(); - console.log(`[test] pool stats after overflow: ${JSON.stringify(stats)}`); - assert( - stats.total <= 1, - `Tracked pool must stay within maxSize after overflow, got total=${stats.total}`, - ); - assert( - server.streamCount() >= CONCURRENCY, - `Expected >= ${CONCURRENCY} streams, got ${server.streamCount()}`, - ); + assertEqual(stats.total, 1, "Pool should stay within maxSize"); + assert(server.streamCount() >= 2, `Expected >= 2 streams, got ${server.streamCount()}`); pool.shutdown(); await server.close(); - console.log("[test] Pool overflow ephemeral workers OK"); + console.log("[test] Pool capacity bound OK"); +} + +async function testProxyMapsPoolCapacityToServiceUnavailable( + modules: TestModules, + backend: TestCursorBackend, +) { + console.log("[test] Testing proxy pool-capacity response..."); + modules.stopProxy(); + backend.setRunMode("text-then-hang"); + const port = await modules.startProxy(async () => "test-token"); + const readers: ReadableStreamDefaultReader[] = []; + + try { + for (let i = 0; i < 4; i++) { + const response = await fetch(`http://localhost:${port}/v1/chat/completions`, { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ + model: "default", + stream: true, + conversation_id: `pool-capacity-${i}`, + messages: [{ role: "user", content: `hold request ${i}` }], + }), + }); + assertEqual(response.status, 200, `Capacity holder ${i} should start`); + assert(response.body, `Capacity holder ${i} should stream`); + const reader = response.body!.getReader(); + readers.push(reader); + await reader.read(); + } + + const rejected = await fetch(`http://localhost:${port}/v1/chat/completions`, { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ + model: "default", + stream: true, + conversation_id: "pool-capacity-rejected", + messages: [{ role: "user", content: "reject while saturated" }], + }), + }); + assertEqual(rejected.status, 503, "Pool saturation should return service unavailable"); + assertEqual(rejected.headers.get("Retry-After"), "2", "Pool saturation should be retryable"); + const payload = await rejected.json() as { error?: { code?: string } }; + assertEqual(payload.error?.code, "service_unavailable", "Pool saturation should use the overload code"); + } finally { + await Promise.all(readers.map((reader) => reader.cancel().catch(() => undefined))); + backend.setRunMode("immediate-close"); + modules.stopProxy(); + } + console.log("[test] Proxy pool-capacity response OK"); } async function testStreamingWatchdogRecoversFromStalledRun( @@ -3053,6 +3108,8 @@ async function main() { // Ephemeral port by default — each process binds an OS-assigned listen port. // OPENCODE_CURSOR_PROXY_PORT remains an optional pin for debugging only. delete process.env.OPENCODE_CURSOR_PROXY_PORT; + delete process.env.OPENCODE_CURSOR_BRIDGE_POOL_MIN; + delete process.env.OPENCODE_CURSOR_BRIDGE_POOL_MAX; const modules = await loadTestModules(); @@ -3076,7 +3133,8 @@ async function main() { await testPersistentBridgeSessionIsolation(); await testPoolRecoveryAfterServerRestart(); await testPoolSequentialRequests(); - await testPoolOverflowEphemeralWorkers(); + await testPoolCapacityBound(); + await testProxyMapsPoolCapacityToServiceUnavailable(modules, backend); await testProxyConsumesCursorModelHeader(modules, backend); await testStreamingWatchdogRecoversFromStalledRun(modules, backend); await testStallExhaustionIsHonest(modules, backend);