t3-code-android-nightly/.repos/alchemy-effect/examples/cloudflare-worker/test/integ.test.ts
Julius Marminge 6f9cea00ae
chore(refs): sync Effect and Alchemy references to 4.0.1 and beta.80 (#16170)
Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
2026-10-05 13:22:30 -07:00

291 lines
11 KiB
TypeScript

import * as Cloudflare from "alchemy/Cloudflare";
import * as Test from "alchemy/Test/Bun";
import { expect } from "bun:test";
import * as Effect from "effect/Effect";
import * as Redacted from "effect/Redacted";
import * as Schedule from "effect/Schedule";
import * as HttpBody from "effect/http/HttpBody";
import * as HttpClient from "effect/http/HttpClient";
import * as HttpClientRequest from "effect/http/HttpClientRequest";
import Stack from "../alchemy.run.ts";
import { WORKFLOW_SECRET_VALUE } from "../src/NotifyWorkflow.ts";
const { test, beforeAll, afterAll, deploy, destroy } = Test.make({
providers: Cloudflare.providers(),
state: Cloudflare.state(),
// dev: true,
});
// This stack deploys a Container (Sandbox) whose image build + push can take
// well over the default 120s hook budget, so give deploy/destroy more room.
const stack = beforeAll(deploy(Stack), { timeout: 600_000 });
afterAll.skipIf(!!process.env.NO_DESTROY)(destroy(Stack), {
timeout: 600_000,
});
// A fresh `workers.dev` URL transiently 404s/5xxs while the route propagates;
// `HttpClient.execute` resolves on those, so retry until the worker answers.
const { executeWhenReady } = Test;
test(
"integ",
Effect.gen(function* () {
const { url } = yield* stack;
expect(url).toBeString();
}),
);
/**
* Better Auth on D1 (auto-migrated at deploy): sign-up + sign-in through
* the `/auth/*` routes served by `auth.fetch`, asserting a session cookie
* comes back. The assets auth panel (`index.html`) drives the same routes
* from the browser.
*/
test(
"better auth: sign-up and sign-in on D1",
Effect.gen(function* () {
const { url } = yield* stack;
const email = "auth-integ@example.com";
const password = "password1234";
const post = (path: string, body: unknown) =>
Effect.tryPromise(async (signal) => {
const response = await fetch(`${url}${path}`, {
method: "POST",
signal,
headers: { "content-type": "application/json" },
body: JSON.stringify(body),
});
return {
status: response.status,
body: await response.text(),
cookies: response.headers.getSetCookie(),
};
});
// A leftover user from a prior NO_DESTROY run is fine — sign-in is the
// real assertion. Retries ride out workers.dev propagation.
yield* post("/auth/sign-up/email", { email, password, name: "Integ" }).pipe(
Effect.filterOrFail(
(r) => r.status === 200 || r.body.includes("USER_ALREADY_EXISTS"),
(r) => new Error(`sign-up failed: ${r.status} ${r.body.slice(0, 200)}`),
),
Effect.retry({ schedule: Schedule.exponential("1 second"), times: 8 }),
);
const signIn = yield* post("/auth/sign-in/email", { email, password });
expect(signIn.status).toBe(200);
expect(signIn.cookies.length).toBeGreaterThan(0);
}),
{ timeout: 120_000 },
);
/**
* Regression guard for https://github.com/alchemy-run/alchemy/pull/172
*
* The stack now includes two Workers (`Api` and `SecondaryApi`) that both
* bind the same `Agent` Durable Object, which in turn binds the `Sandbox`
* Container. Each `yield* Agent` runs the DO's outer init, calling
* `Cloudflare.Container(Sandbox)` once per Worker, so the Sandbox
* ContainerApplication receives two bindings sharing one `namespaceId`.
*
* Before the dedupe fix, `getDurableObjects` counted those as two distinct
* namespaces and the deploy in `beforeAll` died with:
*
* "A Container can only be bound to one Durable Object namespace.
* Found 2 namespaces in bindings: <id>, <id>"
*
* If the deploy ever starts failing again, the whole suite stops at
* `beforeAll` — that is the regression signal. This case just asserts the
* second Worker showed up with a URL so a silent regression that drops the
* binding still surfaces here.
*/
test(
"two workers binding the same container deploy without dedup error",
Effect.gen(function* () {
const { secondaryApiUrl } = yield* stack;
expect(secondaryApiUrl).toBeString();
}),
);
/**
* Regression guard for https://github.com/alchemy-run/alchemy/pull/71
*
* `NotifyWorkflow` accesses `Cloudflare.Workers.WorkerEnvironment` inside its body and
* performs a KV roundtrip via `env.KV.put` / `env.KV.get`. If the fix from #71
* is ever reverted, the body Effect loses the `WorkerEnvironment` service and
* dies with `Service not found: Cloudflare.Workers.WorkerEnvironment` on the
* first `yield* Cloudflare.Workers.WorkerEnvironment` — the workflow instance never
* reaches `complete`, and this test times out or surfaces the `errored` status.
*/
test(
"workflow body can access WorkerEnvironment and exercise env bindings",
Effect.gen(function* () {
const { url } = yield* stack;
interface WorkflowStatus {
status: string;
output?: { secret?: string };
error?: unknown;
}
// Start a fresh workflow instance and poll it to a terminal state. A
// freshly-deployed worker occasionally errors a step while its bindings
// are still propagating, so this effect FAILS on any non-complete
// terminal state, letting the outer retry take another swing with a
// brand-new instance rather than flaking on a transient `errored`.
const runOnce = Effect.gen(function* () {
const roomId = `smoke-${Date.now()}`;
// A freshly-deployed worker transiently 404s (route still propagating)
// or 5xxs (bindings still settling) on the workflow routes, so retry
// through the cold-start window instead of asserting `200` on the first
// hit. Only a non-cold-start status reaches the `expect` below.
const startResponse = yield* executeWhenReady(
HttpClientRequest.post(`${url}/workflow/start/${roomId}`),
);
expect(startResponse.status).toBe(200);
const { instanceId } = (yield* startResponse.json) as {
instanceId: string;
};
expect(instanceId).toBeString();
const client = yield* HttpClient.HttpClient;
const lastStatus = yield* client
.get(`${url}/workflow/status/${instanceId}`)
.pipe(
// Only decode JSON on a 200; a transient 5xx (HTML error page) while
// the worker settles is treated as non-terminal so the poll keeps
// swinging instead of dying on a JSON decode error.
Effect.flatMap((res) =>
res.status === 200
? (res.json as Effect.Effect<unknown, unknown>).pipe(
Effect.map((body) => body as WorkflowStatus),
)
: Effect.succeed({ status: "pending" } as WorkflowStatus),
),
Effect.repeat({
schedule: Schedule.spaced("2 seconds"),
until: (s) => s.status === "complete" || s.status === "errored",
times: 30,
}),
);
// Surface a non-complete terminal state as a failure so the outer retry
// can restart with a fresh instance.
if (lastStatus.status !== "complete") {
return yield* Effect.fail(
new Error(
`workflow ${lastStatus.status}: ${JSON.stringify(lastStatus.error)}`,
),
);
}
return lastStatus;
});
// Cloudflare can briefly route workflow invocations to a worker version
// that predates the final upload (e.g. the pre-create stub), which
// errors instances with "The entrypoint name Notifier was not found in
// this worker" until the deployed version propagates. Each errored
// attempt terminates within a few seconds, so give propagation a
// bounded ~45s of fresh instances rather than 3 swings in 16s.
const lastStatus = yield* runOnce.pipe(
Effect.retry({ schedule: Schedule.spaced("3 seconds"), times: 6 }),
);
expect(lastStatus.status).toBe("complete");
expect(lastStatus.error).toBeFalsy();
// Prove the `Alchemy.Secret(...)` bound at plantime made it all the
// way through to the workflow body's runtime read. The workflow body
// unwraps `Redacted.value(secret)` and embeds it in the returned
// `processed` payload.
expect(lastStatus.output?.secret).toBe(
Redacted.value(WORKFLOW_SECRET_VALUE),
);
}),
{ timeout: 120_000 },
);
/**
* Queue producer→consumer round-trip via the Effect-style
* `Cloudflare.Queues.consumeQueueMessages(Queue, handler)` API.
*
* Producer: `POST /queue/send` returns `{ sent: { id, text, sentAt } }`
* after enqueuing a message.
*
* Consumer: the worker's queue() handler (registered via subscribe in
* src/Api.ts) writes the message body to R2 at `/queue/<id>`. The
* route `GET /queue/result/<id>` reads it back. Cloudflare's queue
* dispatch is async and best-effort, so we poll for up to 60s.
*/
test(
"queue producer→consumer round-trip via consumeQueueMessages()",
Effect.gen(function* () {
const { url } = yield* stack;
const text = `hello-${Date.now()}`;
type Message = { id: string; text: string; sentAt: number };
const send = Effect.gen(function* () {
const sendResponse = yield* executeWhenReady(
HttpClientRequest.post(`${url}/queue/send`).pipe(
HttpClientRequest.setBody(HttpBody.text(text)),
),
);
expect(sendResponse.status).toBe(202);
const { sent } = (yield* sendResponse.json) as { sent: Message };
expect(sent.id).toBeTypeOf("string");
return sent;
});
const sent: Message[] = [yield* send];
// First delivery on a freshly created queue+consumer can lag well past
// the settled steady state (consumer attachment propagates through
// Cloudflare's queue subsystem asynchronously). Poll for ANY of the
// messages we sent, re-sending every ~20s in case an early message was
// published into the propagation window; the consumer persists each
// message to R2 keyed by its id, so the first one to land wins.
let polls = 0;
const findConsumed = Effect.gen(function* () {
polls++;
if (polls % 10 === 0 && sent.length < 5) {
sent.push(yield* send);
}
for (const message of sent) {
const resultResponse = yield* HttpClient.get(
`${url}/queue/result/${message.id}`,
);
if (resultResponse.status === 200) {
return (yield* resultResponse.json) as Message;
}
}
return undefined;
});
const consumed = yield* findConsumed.pipe(
Effect.repeat({
schedule: Schedule.spaced("2 seconds"),
until: (message) => message !== undefined,
times: 75,
}),
);
expect(consumed).toBeDefined();
expect(sent.map((message) => message.id)).toContain(consumed!.id);
expect(consumed!.text).toBe(text);
// Clean up the consumed R2 entries so afterAll's stack.destroy()
// can delete the bucket — otherwise Cloudflare rejects the
// bucket delete with "bucket is not empty".
yield* Effect.forEach(sent, (message) =>
HttpClient.execute(
HttpClientRequest.make("DELETE")(`${url}/queue/result/${message.id}`),
),
);
}),
{ timeout: 180_000 },
);