feat(api): add concurrency limiter to prevent oom under burst load

This commit is contained in:
germondai
2026-06-29 22:11:50 +02:00
parent df79859ebe
commit a42b3a11af
+58 -50
View File
@@ -1,4 +1,4 @@
import { BrowserPool, SessionCache } from "@trawl/browser"
import { BrowserPool, PoolExhaustedError, SessionCache } from "@trawl/browser"
import { scrape } from "@trawl/tiers"
import type { FlareSolverrRequest, FlareSolverrResponse, PoolStats, ScrapeRequest } from "@trawl/types"
import { Elysia } from "elysia"
@@ -47,6 +47,38 @@ function getDeps() {
const startTime = Date.now()
// Concurrency limiter — gates in-flight scrape work at exactly POOL_SIZE so we
// reject immediately (HTTP 429) instead of letting BrowserPool.acquire() block
// for up to 5s. Prowlarr/Jackett use the 429 to back off cleanly.
let activeJobs = 0
const MAX_CONCURRENT = POOL_SIZE
function tryAcquireSlot(): boolean {
if (activeJobs < MAX_CONCURRENT) {
activeJobs++
return true
}
return false
}
function releaseSlot(): void {
if (activeJobs > 0) activeJobs--
}
// Build a FlareSolverr v2-shaped error envelope. Used by /v1, /scrape, and the
// second-line PoolExhaustedError catch-arm.
function flareSolverrError(url: string, message: string): FlareSolverrResponse {
const now = Date.now()
return {
status: "error",
message,
startTimestamp: now,
endTimestamp: now,
version: "2.0.0",
solution: { url, status: 0, headers: {}, response: "", cookies: [], userAgent: "" },
}
}
new Elysia()
.get("/health", () => ({
status: pool ? "ok" : "starting",
@@ -76,6 +108,8 @@ new Elysia()
busy: stats.busy,
restarts: stats.restarts,
queueDepth: 0,
activeJobs,
maxConcurrent: MAX_CONCURRENT,
}
})
@@ -87,40 +121,17 @@ new Elysia()
if (cmd !== "request.get" && cmd !== "request.post") {
set.status = 400
return {
status: "error",
message: `Unknown cmd: ${cmd}`,
startTimestamp,
endTimestamp: Date.now(),
version: "2.0.0",
solution: {
url: "",
status: 0,
headers: {},
response: "",
cookies: [],
userAgent: "",
},
} satisfies FlareSolverrResponse
return flareSolverrError(req.url, `Unknown cmd: ${cmd}`)
}
if (!pool) {
set.status = 503
return {
status: "error",
message: "Browser pool initializing, retry in a few seconds",
startTimestamp,
endTimestamp: Date.now(),
version: "2.0.0",
solution: {
url: "",
status: 0,
headers: {},
response: "",
cookies: [],
userAgent: "",
},
} satisfies FlareSolverrResponse
return flareSolverrError(req.url, "Browser pool initializing, retry in a few seconds")
}
if (!tryAcquireSlot()) {
set.status = 429
return flareSolverrError(req.url, "Browser pool saturated, retry shortly")
}
try {
@@ -141,22 +152,10 @@ new Elysia()
},
} satisfies FlareSolverrResponse
} catch (err) {
set.status = 500
return {
status: "error",
message: err instanceof Error ? err.message : String(err),
startTimestamp,
endTimestamp: Date.now(),
version: "2.0.0",
solution: {
url: req.url,
status: 0,
headers: {},
response: "",
cookies: [],
userAgent: "",
},
} satisfies FlareSolverrResponse
set.status = err instanceof PoolExhaustedError ? 429 : 500
return flareSolverrError(req.url, err instanceof Error ? err.message : String(err))
} finally {
releaseSlot()
}
})
@@ -164,14 +163,23 @@ new Elysia()
.post("/scrape", async ({ body, set }) => {
if (!pool) {
set.status = 503
return { error: "Browser pool initializing, retry in a few seconds" }
return flareSolverrError("", "Browser pool initializing, retry in a few seconds")
}
const req = body as ScrapeRequest
if (!tryAcquireSlot()) {
set.status = 429
return flareSolverrError(req.url ?? "", "Browser pool saturated, retry shortly")
}
try {
return await scrape(req, getDeps())
} catch (err) {
set.status = 500
return { error: err instanceof Error ? err.message : String(err) }
set.status = err instanceof PoolExhaustedError ? 429 : 500
return flareSolverrError(req.url ?? "", err instanceof Error ? err.message : String(err))
} finally {
releaseSlot()
}
})