import { Credentials, oauthCredentials, Retry, type CloudflareOpContext, } from "@distilled.cloud/cloudflare"; import * as Effect from "effect/Effect"; import * as FetchHttpClient from "effect/unstable/http/FetchHttpClient"; import { bodyBytes, withResponse } from "../server/http.ts"; // No process credentials, refresh cache or nested retries. The calling request // or native Workflow attempt has already resolved its own authorized grant. export function cloudflare( operation: Effect.Effect, accessToken: string, network: typeof fetch, timeout = 60_000, limit = 2 * 1024 * 1024, ) { const boundedFetch: typeof fetch = (input, init) => Effect.runPromise( withResponse( network, input, init ?? {}, timeout, new Error("Cloudflare transport unavailable"), (response) => Effect.gen(function* () { if (response.status >= 300 && response.status < 400) return yield* Effect.die( new Error("Cloudflare redirect rejected"), ); if (!response.body) return response; const bytes = yield* bodyBytes( response, limit, new Error("Cloudflare response unavailable"), ); // The SDK also accepts bare JSON for other services. These control-plane // operations require Cloudflare's successful v4 envelope explicitly. if (response.ok) { yield* Effect.try({ try: () => { const envelope: unknown = JSON.parse( new TextDecoder().decode(bytes), ); if ( !envelope || typeof envelope !== "object" || !("success" in envelope) || envelope.success !== true || !("result" in envelope) ) throw new Error("Invalid Cloudflare envelope"); }, catch: () => new Error("Cloudflare response unavailable"), }); } return new Response(bytes, response); }), ), { signal: init?.signal ?? undefined }, ); return operation.pipe( Retry.none, Effect.provideService( Credentials, Effect.succeed(oauthCredentials({ accessToken })), ), Effect.provide(FetchHttpClient.layer), Effect.provideService(FetchHttpClient.Fetch, boundedFetch), Effect.provideService(FetchHttpClient.RequestInit, { redirect: "manual" }), ); }