| 1 | export class BodyReadError extends Error { |
| 2 | constructor(readonly status: 400 | 413, message: string) { |
| 3 | super(message); |
| 4 | this.name = "BodyReadError"; |
| 5 | } |
| 6 | } |
| 7 | |
| 8 | /** Count raw bytes while reading; Content-Length is only an early rejection. */ |
| 9 | export async function readBoundedBody(request: { body: ReadableStream<Uint8Array> | null; headers: Headers }, maxBytes: number): Promise<Uint8Array> { |
| 10 | const reject = async (status: 400 | 413, message: string): Promise<never> => { |
| 11 | try { await request.body?.cancel(message); } catch { /* Keep the original rejection. */ } |
| 12 | throw new BodyReadError(status, message); |
| 13 | }; |
| 14 | const rawLength = request.headers.get("content-length"); |
| 15 | if (rawLength !== null) { |
| 16 | if (!/^\d+$/.test(rawLength)) return reject(400, "invalid Content-Length"); |
| 17 | if (Number(rawLength) > maxBytes) return reject(413, "payload too large"); |
| 18 | } |
| 19 | if (!request.body) return new Uint8Array(); |
| 20 | |
| 21 | const reader = request.body.getReader(); |
| 22 | const chunks: Uint8Array[] = []; |
| 23 | let total = 0; |
| 24 | try { |
| 25 | while (true) { |
| 26 | const { done, value } = await reader.read(); |
| 27 | if (done) break; |
| 28 | total += value.byteLength; |
| 29 | if (total > maxBytes) throw new BodyReadError(413, "payload too large"); |
| 30 | chunks.push(value); |
| 31 | } |
| 32 | } catch (cause) { |
| 33 | try { await reader.cancel("body rejected"); } catch { /* Preserve the read/size error. */ } |
| 34 | throw cause instanceof BodyReadError ? cause : new BodyReadError(400, "body read failed"); |
| 35 | } finally { |
| 36 | reader.releaseLock(); |
| 37 | } |
| 38 | const bytes = new Uint8Array(total); |
| 39 | let offset = 0; |
| 40 | for (const chunk of chunks) { |
| 41 | bytes.set(chunk, offset); |
| 42 | offset += chunk.byteLength; |
| 43 | } |
| 44 | return bytes; |
| 45 | } |
| 46 | |
| 47 | /** Budget for background (cron) reads of other hosts. */ |
| 48 | export interface OutboundBudget { |
| 49 | /** Deadline for the whole exchange, body included. */ |
| 50 | timeoutMs?: number; |
| 51 | /** Largest body read; a larger one is an error, never a silent truncation. */ |
| 52 | maxBytes?: number; |
| 53 | } |
| 54 | |
| 55 | export const OUTBOUND_TIMEOUT_MS = 15_000; |
| 56 | export const OUTBOUND_MAX_BYTES = 2 * 1024 * 1024; |
| 57 | |
| 58 | /** |
| 59 | * One outbound GET-style read under a deadline and a byte cap. A slow or |
| 60 | * endless peer cannot hold a cron invocation until the platform kills it, and |
| 61 | * a huge body cannot exhaust isolate memory. Non-2xx answers return |
| 62 | * `{ ok: false }` with the body discarded; transport, timeout and size |
| 63 | * failures throw, so callers keep their own fallbacks. |
| 64 | */ |
| 65 | export async function fetchBoundedText( |
| 66 | url: string, |
| 67 | init: RequestInit = {}, |
| 68 | { timeoutMs = OUTBOUND_TIMEOUT_MS, maxBytes = OUTBOUND_MAX_BYTES }: OutboundBudget = {}, |
| 69 | ): Promise<{ ok: boolean; status: number; text: string }> { |
| 70 | const response = await fetch(url, { ...init, signal: AbortSignal.timeout(timeoutMs) }); |
| 71 | if (!response.ok) { |
| 72 | try { await response.body?.cancel(); } catch { /* The status is the answer. */ } |
| 73 | return { ok: false, status: response.status, text: "" }; |
| 74 | } |
| 75 | const bytes = await readBoundedBody(response, maxBytes); |
| 76 | return { ok: true, status: response.status, text: new TextDecoder().decode(bytes) }; |
| 77 | } |
| 78 |