Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20 changes: 19 additions & 1 deletion packages/alchemy/src/Cloudflare/Hyperdrive/ConnectBinding.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ import * as Effect from "effect/Effect";
import * as Layer from "effect/Layer";
import * as Redacted from "effect/Redacted";
import * as Output from "../../Output.ts";
import { defaultProviderMode, type ProviderMode } from "../../ProviderMode.ts";
import { Worker, WorkerEnvironment } from "../Workers/Worker.ts";
import { Connect, type ConnectClient } from "./Connect.ts";
import type { Connection } from "./Connection.ts";
Expand All @@ -24,7 +25,7 @@ export const ConnectBinding = Layer.effect(
id: connection.hyperdriveId as unknown as string,
},
],
hyperdrives: getHyperdriveDevOrigin(connection),
hyperdrives: yield* getHyperdriveDevOriginForHost(connection, host),
});
}

Expand All @@ -45,6 +46,23 @@ export const ConnectBinding = Layer.effect(
}),
);

/**
* The `hyperdrives` dev channel for `connection` bound to `host`. Only the
* local worker provider reads it, so it is contributed only when the host
* runs locally — its registration-captured `Mode` (`Alchemy.remote()` →
* `"live"`) or the run default (`alchemy dev` → `"local"`), the same
* resolution the planner applies to the host. A live host (`alchemy deploy`,
* or a `remote()` worker in dev) never evaluates the dev origin, so an
* Access-protected origin without a `dev` override deploys.
*/
export const getHyperdriveDevOriginForHost = Effect.fn(function* (
connection: Connection,
host: { readonly Mode?: ProviderMode | undefined },
) {
const mode = host.Mode ?? (yield* defaultProviderMode);
return mode === "local" ? getHyperdriveDevOrigin(connection) : undefined;
});

export const getHyperdriveDevOrigin = (connection: Connection) => {
const origin = Output.map(
Output.all(connection.dev, connection.origin, connection.mtls),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@ import type { ContainerApplication } from "../Containers/ContainerApplication.ts
import { isDatabase } from "../D1/Database.ts";
import { isSendEmail } from "../Email/SendEmail.ts";
import { isApp } from "../Flagship/App.ts";
import { getHyperdriveDevOrigin } from "../Hyperdrive/ConnectBinding.ts";
import { getHyperdriveDevOriginForHost } from "../Hyperdrive/ConnectBinding.ts";
import { isHyperdriveConnection } from "../Hyperdrive/Connection.ts";
import { isImages } from "../Images/Images.ts";
import { isNamespace as isKVNamespace } from "../KV/Namespace.ts";
Expand Down Expand Up @@ -263,7 +263,7 @@ export const bindWorkerAsyncBindings = Effect.fn(function* (
yield* resource.bind`${bindingName}`({
bindings: [resolvedBindingMeta],
hyperdrives: isHyperdriveConnection(binding)
? getHyperdriveDevOrigin(binding)
? yield* getHyperdriveDevOriginForHost(binding, resource)
: undefined,
// Dev-only local-emulation opt-out channel (like `hyperdrives`):
// worker-only bindings and `SendEmail` descriptors piped through
Expand Down
202 changes: 202 additions & 0 deletions packages/alchemy/test/Cloudflare/Hyperdrive/Hyperdrive.test.ts
Original file line number Diff line number Diff line change
@@ -1,15 +1,30 @@
import * as hyperdrive from "@distilled.cloud/cloudflare/hyperdrive";
import * as zeroTrust from "@distilled.cloud/cloudflare/zero-trust";
import { assert, expect } from "alchemy-test";
import * as Data from "effect/Data";
import * as Effect from "effect/Effect";
import * as HttpClient from "effect/http/HttpClient";
import * as Layer from "effect/Layer";
import * as Redacted from "effect/Redacted";
import { MinimumLogLevel } from "effect/References";
import * as Schedule from "effect/Schedule";
import * as pathe from "pathe";
import * as Cloudflare from "@/Cloudflare";
import { CloudflareEnvironment } from "@/Cloudflare/CloudflareEnvironment";
import { findZoneByName } from "@/Cloudflare/Zone/lookup";
import * as Neon from "@/Neon";
import * as Provider from "@/Provider";
import * as Test from "@/Test/Alchemy";
import { waitForWorkerToBeDeleted } from "../Utils/Worker.ts";
import HyperdriveAccessEffectWorker from "./fixtures/access-effect-worker.ts";
import {
ACCESS_ORIGIN_HOST,
ACCESS_ORIGIN_PASSWORD,
ACCESS_ORIGIN_PORT,
ACCESS_ORIGIN_ZONE,
AccessOriginConnection,
AccessOriginRoute,
} from "./fixtures/access-origin.ts";

const { test } = Test.make({ providers: Layer.merge(Cloudflare.providers(), Neon.providers()) });

Expand Down Expand Up @@ -137,6 +152,193 @@ test.provider(
},
);

// ── Access-protected origin (#1836) ────────────────────────────────────────
// A real origin behind Cloudflare Access: Postgres in Docker, published
// through a Cloudflare Tunnel whose connector (`cloudflared`) runs here,
// guarded by a self-hosted Access app that admits one service token. Both
// binding flavors bind the Connection with no `dev` override and must
// deploy and query through it.

const cloudflaredBin = Bun.which("cloudflared");
const dockerBin = Bun.which("docker");
const ACCESS_ORIGIN_CONTAINER = "alchemy-test-hyperdrive-access-origin";

const run = (cmd: string[]) =>
Effect.sync(() => Bun.spawnSync(cmd, { stdout: "pipe", stderr: "pipe" }));

/** Postgres with TLS (the image's snakeoil cert) on {@link ACCESS_ORIGIN_PORT}. */
const accessOriginPostgres = Effect.acquireRelease(
Effect.gen(function* () {
yield* run([dockerBin!, "rm", "-f", ACCESS_ORIGIN_CONTAINER]);
const started = yield* run([
dockerBin!,
"run",
"-d",
"--rm",
"--name",
ACCESS_ORIGIN_CONTAINER,
"-e",
`POSTGRES_PASSWORD=${ACCESS_ORIGIN_PASSWORD}`,
"-p",
`127.0.0.1:${ACCESS_ORIGIN_PORT}:5432`,
"postgres:17",
"-c",
"ssl=on",
"-c",
"ssl_cert_file=/etc/ssl/certs/ssl-cert-snakeoil.pem",
"-c",
"ssl_key_file=/etc/ssl/private/ssl-cert-snakeoil.key",
]);
if (started.exitCode !== 0) {
return yield* Effect.die(new Error(`docker run failed: ${started.stderr.toString()}`));
}
yield* run([
dockerBin!,
"exec",
ACCESS_ORIGIN_CONTAINER,
"pg_isready",
"-h",
"127.0.0.1",
"-U",
"postgres",
]).pipe(
Effect.repeat({
schedule: Schedule.spaced("1 second"),
until: (result) => result.exitCode === 0,
times: 30,
}),
);
}),
() => run([dockerBin!, "rm", "-f", ACCESS_ORIGIN_CONTAINER]),
);

class TunnelNotHealthy extends Data.TaggedError("TunnelNotHealthy")<{ status: string }> {}
class AccessQueryFailed extends Data.TaggedError("AccessQueryFailed")<{
status: number;
body: string;
}> {}

/** Run the tunnel's connector here until the scope closes, then wait for it to register. */
const accessOriginConnector = (accountId: string, tunnelId: string, token: string) =>
Effect.gen(function* () {
yield* Effect.acquireRelease(
Effect.sync(() =>
Bun.spawn([cloudflaredBin!, "tunnel", "run", "--token", token], {
stdout: "ignore",
stderr: "ignore",
}),
),
(proc) =>
Effect.promise(async () => {
proc.kill();
await proc.exited;
}),
);
yield* zeroTrust.getTunnelCloudflared({ accountId, tunnelId }).pipe(
Effect.flatMap((tunnel) =>
tunnel.status === "healthy"
? Effect.void
: Effect.fail(new TunnelNotHealthy({ status: String(tunnel.status) })),
),
Effect.retry({
while: (e) => e._tag === "TunnelNotHealthy",
schedule: Schedule.spaced("2 seconds"),
times: 30,
}),
);
});

const queryThroughAccess = (url: string) =>
HttpClient.get(url).pipe(
Effect.flatMap((res) =>
res.text.pipe(
Effect.flatMap((body) =>
res.status === 200
? Effect.succeed(JSON.parse(body) as { via: string })
: Effect.fail(new AccessQueryFailed({ status: res.status, body })),
),
),
),
Effect.retry({ schedule: Schedule.spaced("3 seconds"), times: 20 }),
);

test.provider.skipIf(!cloudflaredBin || !dockerBin)(
"deploys and queries an Access-protected origin without a dev override",
(stack) =>
Effect.gen(function* () {
const { accountId } = yield* yield* CloudflareEnvironment;
const zone = yield* findZoneByName({ accountId, name: ACCESS_ORIGIN_ZONE });
if (!zone) return yield* Effect.die(new Error(`zone ${ACCESS_ORIGIN_ZONE} not found`));

yield* stack.destroy();

const deployed = yield* Effect.scoped(
Effect.gen(function* () {
yield* accessOriginPostgres;

// The tunnel must have a live connector before Hyperdrive is
// created: Cloudflare connects to the origin when it creates it.
const route = yield* stack.deploy(AccessOriginRoute(zone.id));
yield* accessOriginConnector(
accountId,
route.tunnel.tunnelId,
Redacted.value(route.tunnel.token),
);

const deployed = yield* stack
.deploy(
Effect.gen(function* () {
yield* AccessOriginRoute(zone.id);
const connection = yield* AccessOriginConnection;
const effectWorker = yield* HyperdriveAccessEffectWorker;
const asyncWorker = yield* Cloudflare.Worker("HyperdriveAccessAsyncWorker", {
main: pathe.resolve(import.meta.dirname, "fixtures/access-async-worker.ts"),
env: { HD: AccessOriginConnection },
});
return { connection, effectWorker, asyncWorker };
}),
)
.pipe(
// Hyperdrive resolves the origin host when it creates the config,
// and a just-created CNAME can take a little while to answer.
Effect.retry({ schedule: Schedule.spaced("10 seconds"), times: 12 }),
);

const actual = yield* hyperdrive.getConfig({
accountId,
hyperdriveId: deployed.connection.hyperdriveId,
});
assert("accessClientId" in actual.origin, "origin must be Access-protected");
expect(actual.origin.host).toEqual(ACCESS_ORIGIN_HOST);

// Both binding flavors reach Postgres through Access + the tunnel.
expect(yield* queryThroughAccess(deployed.effectWorker.url!)).toEqual({
via: "through-access",
});
expect(yield* queryThroughAccess(deployed.asyncWorker.url!)).toEqual({
via: "through-access",
});
return deployed;
}),
);

yield* stack.destroy();
yield* waitForConfigToBeDeleted(deployed.connection.hyperdriveId, accountId);
yield* waitForWorkerToBeDeleted(deployed.effectWorker.workerName, accountId);
yield* waitForWorkerToBeDeleted(deployed.asyncWorker.workerName, accountId);
}).pipe(logLevel),
{
tags: [
"provider:cloudflare",
"provider:cloudflare:hyperdrive",
"provider:cloudflare:tunnel",
"provider:cloudflare:access",
"live",
],
timeout: 300_000,
},
);

const waitForConfigToBeDeleted = Effect.fn(function* (hyperdriveId: string, accountId: string) {
yield* hyperdrive.getConfig({ accountId, hyperdriveId }).pipe(
Effect.flatMap(() => Effect.fail(new ConfigStillExists())),
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,21 @@
import type { Hyperdrive } from "@cloudflare/workers-types";
import { Client } from "pg";

/**
* Async Worker that binds the Access-protected Hyperdrive via
* `env: { HD: connection }` and runs a query through it.
*/
export default {
async fetch(_request: Request, env: { HD: Hyperdrive }): Promise<Response> {
const client = new Client({ connectionString: env.HD.connectionString });
try {
await client.connect();
const result = await client.query("select 'through-access' as via");
return Response.json(result.rows[0]);
} catch (error) {
return new Response(String(error), { status: 500 });
} finally {
await client.end().catch(() => {});
}
},
};
Original file line number Diff line number Diff line change
@@ -0,0 +1,40 @@
import * as Effect from "effect/Effect";
import * as HttpServerResponse from "effect/http/HttpServerResponse";
import * as Redacted from "effect/Redacted";
import { Client } from "pg";
import * as Cloudflare from "@/Cloudflare/index.ts";
import { AccessOriginConnection } from "./access-origin.ts";

/**
* Effect Worker that binds the Access-protected Hyperdrive through
* `Cloudflare.Hyperdrive.Connect` and runs a query through it, so a 200
* proves the whole path: Hyperdrive → Access → Tunnel → Postgres.
*/
export default class HyperdriveAccessEffectWorker extends Cloudflare.Worker<HyperdriveAccessEffectWorker>()(
"HyperdriveAccessEffectWorker",
{ main: import.meta.url },
Effect.gen(function* () {
const connection = yield* AccessOriginConnection;
const hd = yield* Cloudflare.Hyperdrive.Connect(connection);
return {
fetch: Effect.gen(function* () {
const connectionString = Redacted.value(yield* hd.connectionString);
return yield* Effect.tryPromise(async () => {
const client = new Client({ connectionString });
await client.connect();
try {
const result = await client.query("select 'through-access' as via");
return result.rows[0] as { via: string };
} finally {
await client.end();
}
}).pipe(
Effect.flatMap((row) => HttpServerResponse.json(row)),
Effect.catch((error) =>
Effect.succeed(HttpServerResponse.text(String(error), { status: 500 })),
),
);
}),
};
}).pipe(Effect.provide(Cloudflare.Hyperdrive.ConnectBinding)),
) {}
Loading
Loading