From 96e53df180cd947a724e664a26ac295f67ada3a1 Mon Sep 17 00:00:00 2001 From: Ujjwal Reddy Date: Wed, 12 Aug 2026 15:02:30 -0400 Subject: [PATCH 1/3] fix(js-sdk): enforce inflight concurrency cap on streaming bodies --- packages/js-sdk/src/api/inflight.ts | 77 ++++++++++++++++++++++++++++- 1 file changed, 75 insertions(+), 2 deletions(-) diff --git a/packages/js-sdk/src/api/inflight.ts b/packages/js-sdk/src/api/inflight.ts index 16401467af..9951959ca1 100644 --- a/packages/js-sdk/src/api/inflight.ts +++ b/packages/js-sdk/src/api/inflight.ts @@ -73,10 +73,83 @@ export function limitConcurrency( const signal = init?.signal ?? (isRequestLike(input) ? input.signal : undefined) const release = await sem.acquire(signal) + let res: Response try { - return await fetcher(input, init) - } finally { + res = await fetcher(input, init) + } catch (err) { release() + throw err } + + if (!res.body) { + release() + return res + } + + let released = false + const safeRelease = () => { + if (!released) { + released = true + release() + } + } + + return createResponseProxy(res, safeRelease) }) as typeof fetch } + +function createResponseProxy(res: Response, safeRelease: () => void): Response { + const methods = ['arrayBuffer', 'blob', 'formData', 'json', 'text'] as const + for (const method of methods) { + const original = (res[method] as any).bind(res) + ;(res as any)[method] = async (...args: any[]) => { + try { + return await original(...args) + } finally { + safeRelease() + } + } + } + + if (!res.body) return res + + const bodyProxy = new Proxy(res.body as any, { + get(target, prop, receiver) { + if (prop === 'getReader') { + return function (...args: any[]) { + const reader = target.getReader(...args) + const originalRead = reader.read.bind(reader) + const originalCancel = reader.cancel.bind(reader) + + reader.read = async () => { + try { + const result = await originalRead() + if (result.done) safeRelease() + return result + } catch (err) { + safeRelease() + throw err + } + } + + reader.cancel = async (reason?: any) => { + try { + return await originalCancel(reason) + } finally { + safeRelease() + } + } + return reader + } + } + return Reflect.get(target, prop, receiver) + } + }) + + return new Proxy(res, { + get(target, prop, receiver) { + if (prop === 'body') return bodyProxy + return Reflect.get(target, prop, receiver) + } + }) +} From 8221490b5fe8e216aacc41510a604e9a60bbfed4 Mon Sep 17 00:00:00 2001 From: Ujjwal Reddy Date: Wed, 12 Aug 2026 15:07:49 -0400 Subject: [PATCH 2/3] chore: add changeset for inflight concurrency fix --- .changeset/fix-inflight-concurrency.md | 7 +++++++ 1 file changed, 7 insertions(+) create mode 100644 .changeset/fix-inflight-concurrency.md diff --git a/.changeset/fix-inflight-concurrency.md b/.changeset/fix-inflight-concurrency.md new file mode 100644 index 0000000000..f7cd909960 --- /dev/null +++ b/.changeset/fix-inflight-concurrency.md @@ -0,0 +1,7 @@ +--- +"e2b": patch +--- + +fix(js-sdk): enforce inflight concurrency cap on streaming bodies + +Defer the release of the inflight semaphore until the `Response.body` is fully consumed, aborted, or errors out. This prevents HTTP/2 stream exhaustion under heavy load. From 6442b84e399327c3f94cff75a0f3ab59b11e09a1 Mon Sep 17 00:00:00 2001 From: Ujjwal Reddy Date: Wed, 12 Aug 2026 15:12:34 -0400 Subject: [PATCH 3/3] fix(js-sdk): address codex review feedback on inflight concurrency proxy --- packages/js-sdk/src/api/inflight.ts | 85 +++++++++++++++++++++-------- 1 file changed, 63 insertions(+), 22 deletions(-) diff --git a/packages/js-sdk/src/api/inflight.ts b/packages/js-sdk/src/api/inflight.ts index 9951959ca1..172899074c 100644 --- a/packages/js-sdk/src/api/inflight.ts +++ b/packages/js-sdk/src/api/inflight.ts @@ -113,43 +113,84 @@ function createResponseProxy(res: Response, safeRelease: () => void): Response { if (!res.body) return res - const bodyProxy = new Proxy(res.body as any, { - get(target, prop, receiver) { - if (prop === 'getReader') { - return function (...args: any[]) { - const reader = target.getReader(...args) - const originalRead = reader.read.bind(reader) - const originalCancel = reader.cancel.bind(reader) - - reader.read = async () => { + const createBodyProxy = (stream: any): any => { + return new Proxy(stream, { + get(target, prop) { + if (prop === 'cancel') { + return async (...args: any[]) => { try { - const result = await originalRead() - if (result.done) safeRelease() - return result - } catch (err) { + return await target.cancel(...args) + } finally { + safeRelease() + } + } + } + if (prop === 'pipeTo') { + return async (...args: any[]) => { + try { + return await target.pipeTo(...args) + } finally { safeRelease() - throw err } } + } + if (prop === 'pipeThrough') { + return (...args: any[]) => { + const newStream = target.pipeThrough(...args) + return createBodyProxy(newStream) + } + } + if (prop === 'getReader') { + return function (...args: any[]) { + const reader = target.getReader(...args) + const originalRead = reader.read.bind(reader) + const originalCancel = reader.cancel.bind(reader) + + reader.read = async (...readArgs: any[]) => { + try { + const result = await originalRead(...readArgs) + if (result.done) safeRelease() + return result + } catch (err) { + safeRelease() + throw err + } + } - reader.cancel = async (reason?: any) => { + reader.cancel = async (reason?: any) => { + try { + return await originalCancel(reason) + } finally { + safeRelease() + } + } + return reader + } + } + if (prop === Symbol.asyncIterator) { + if (!target[Symbol.asyncIterator]) return undefined + return async function* (...args: any[]) { try { - return await originalCancel(reason) + yield* target[Symbol.asyncIterator](...args) } finally { safeRelease() } } - return reader } + + const value = Reflect.get(target, prop) + return typeof value === 'function' ? value.bind(target) : value } - return Reflect.get(target, prop, receiver) - } - }) + }) + } + + const bodyProxy = createBodyProxy(res.body) return new Proxy(res, { - get(target, prop, receiver) { + get(target, prop) { if (prop === 'body') return bodyProxy - return Reflect.get(target, prop, receiver) + const value = Reflect.get(target, prop) + return typeof value === 'function' ? value.bind(target) : value } }) }