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,
}),
),
);
}