Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions .changeset/fix-inflight-concurrency.md
Original file line number Diff line number Diff line change
@@ -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.
118 changes: 116 additions & 2 deletions packages/js-sdk/src/api/inflight.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
})
}