import * as Effect from "effect/Effect"; export function withAbortSignal( use: (signal: AbortSignal) => Effect.Effect, ) { return Effect.acquireUseRelease( Effect.sync(() => new AbortController()), (controller) => use(controller.signal), (controller) => Effect.sync(() => controller.abort()), ); } // A request owns its body as well as its headers. Native cancellation and the // Effect deadline cover slow body reads and release the connection on failure. export function withResponse( network: typeof fetch, input: string | URL | Request, init: RequestInit, timeout: number, unavailable: E, read: (response: Response) => Effect.Effect, ) { return Effect.acquireUseRelease( Effect.sync(() => new AbortController()), (controller) => Effect.tryPromise({ try: (signal) => network(input, { ...init, redirect: "manual", signal: AbortSignal.any([ signal, controller.signal, ...(init.signal ? [init.signal] : []), ...(input instanceof Request ? [input.signal] : []), ]), }), catch: () => unavailable, }).pipe( Effect.flatMap(read), Effect.timeoutOrElse({ duration: timeout, orElse: () => Effect.fail(unavailable), }), ), (controller) => Effect.sync(() => controller.abort()), ); } export function bodyBytes( response: Response, limit: number, invalid: E, unavailable: E | F = invalid, ) { return Effect.gen(function* () { const body = response.body; if (!body) return yield* Effect.fail(invalid); return yield* Effect.acquireUseRelease( Effect.sync(() => body.getReader()), (reader) => Effect.gen(function* () { if (Number(response.headers.get("Content-Length")) > limit) return yield* Effect.fail(invalid); const chunks: Uint8Array[] = []; let size = 0; for (;;) { const chunk = yield* Effect.tryPromise({ try: () => reader.read(), catch: () => unavailable, }); if (chunk.done) break; size += chunk.value.byteLength; if (size > limit) return yield* Effect.fail(invalid); chunks.push(chunk.value); } const bytes = new Uint8Array(size); let offset = 0; for (const chunk of chunks) { bytes.set(chunk, offset); offset += chunk.byteLength; } return bytes; }), (reader) => Effect.tryPromise({ try: () => reader.cancel(), catch: () => undefined, }).pipe( Effect.timeoutOrElse({ duration: 1_000, orElse: () => Effect.void }), Effect.ensuring(Effect.sync(() => reader.releaseLock())), Effect.ignore, ), ); }); } export function bodyJson(response: Response, limit: number, invalid: E) { return bodyBytes(response, limit, invalid).pipe( Effect.flatMap((bytes) => Effect.try({ try: (): unknown => JSON.parse(new TextDecoder().decode(bytes)), catch: () => invalid, }), ), ); }