Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108import * as Effect from "effect/Effect";
export function withAbortSignal<A, E, R>( use: (signal: AbortSignal) => Effect.Effect<A, E, R>,) { 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<A, E, F>( network: typeof fetch, input: string | URL | Request, init: RequestInit, timeout: number, unavailable: E, read: (response: Response) => Effect.Effect<A, F>,) { 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<E, F = never>( 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<E>(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, }), ), );}