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. diff --git a/packages/js-sdk/src/api/inflight.ts b/packages/js-sdk/src/api/inflight.ts index 16401467af..172899074c 100644 --- a/packages/js-sdk/src/api/inflight.ts +++ b/packages/js-sdk/src/api/inflight.ts @@ -73,10 +73,124 @@ 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 createBodyProxy = (stream: any): any => { + return new Proxy(stream, { + get(target, prop) { + if (prop === 'cancel') { + return async (...args: any[]) => { + try { + return await target.cancel(...args) + } finally { + safeRelease() + } + } + } + if (prop === 'pipeTo') { + return async (...args: any[]) => { + try { + return await target.pipeTo(...args) + } finally { + safeRelease() + } + } + } + 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) => { + 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 { + yield* target[Symbol.asyncIterator](...args) + } finally { + safeRelease() + } + } + } + + const value = Reflect.get(target, prop) + return typeof value === 'function' ? value.bind(target) : value + } + }) + } + + const bodyProxy = createBodyProxy(res.body) + + return new Proxy(res, { + get(target, prop) { + if (prop === 'body') return bodyProxy + const value = Reflect.get(target, prop) + return typeof value === 'function' ? value.bind(target) : value + } + }) +}