import assert from "node:assert/strict"; import { test } from "node:test"; import { build } from "esbuild"; const bundle = await build({ stdin: { contents: ` export * as Effect from 'effect/Effect'; export { searchWeb } from './worker/web-search.ts'; export { runWeb, WebFailure } from './worker/web-errors.ts'; export { searchResponse } from './tests/fixtures/search-response.ts'; `, resolveDir: process.cwd(), }, bundle: true, platform: "node", format: "esm", write: false, }); const { Effect, searchWeb, runWeb, WebFailure, searchResponse } = await import( `data:text/javascript;base64,${Buffer.from(bundle.outputFiles[0].text).toString("base64")}` ); test("search programs are lazy and keep invocation state and sources isolated", async () => { const requests = []; const ai = { async run(model, input, options) { requests.push({ model, input, options }); return Response.json( searchResponse( `https://${input.input}.example.com/page`, input.input, "Evidence", ), ); }, }; const first = searchWeb(ai, "first", 3); const second = searchWeb(ai, "second", 3); assert.equal(requests.length, 0); const values = await Promise.all([ runWeb(first), runWeb(second), runWeb(first), ]); assert.deepEqual( values.map((value) => value.sources[0].title), ["first", "second", "first"], ); assert.equal(values[0].sources[0].id, values[2].sources[0].id); assert.notEqual(requests[0].options.signal, requests[2].options.signal); assert.ok(requests.every((request) => request.options.signal.aborted)); }); test("cancellation before execution makes no native request", async () => { let calls = 0; const ai = { async run() { calls++; return Response.json({}); }, }; const signal = AbortSignal.abort(); assert.equal( (await runWeb(searchWeb(ai, "query", 3), signal)).code, "cancelled", ); assert.equal(calls, 0); }); test( "interrupting a search body cancels its reader and native binding signal", { timeout: 2000 }, async () => { const reading = Promise.withResolvers(); let cancelled = 0; let bindingSignal; let body; const ai = { async run(_model, _input, options) { bindingSignal = options.signal; body = new ReadableStream({ pull() { reading.resolve(); }, cancel() { cancelled++; }, }); return new Response(body); }, }; const caller = new AbortController(); const result = runWeb(searchWeb(ai, "query", 3), caller.signal); await reading.promise; caller.abort(); assert.equal((await result).code, "cancelled"); assert.equal(cancelled, 1); assert.equal(body.locked, false); assert.equal(bindingSignal.aborted, true); }, ); test( "deadline interrupts an unresponsive binding and retains the timeout failure", { timeout: 2000 }, async () => { let signal; const ai = { run(_model, _input, options) { signal = options.signal; return new Promise(() => {}); }, }; const program = searchWeb(ai, "query", 3).pipe( Effect.timeoutOrElse({ duration: 10, orElse: () => Effect.fail(new WebFailure("timeout")), }), ); assert.equal((await runWeb(program)).code, "timeout"); assert.equal(signal.aborted, true); }, ); test("binding failures and defects never expose provider payloads in tool results", async () => { const secret = "private-provider-response-and-credentials"; const ai = { async run() { throw new Error(secret); }, }; const unavailable = await runWeb(searchWeb(ai, "query", 3)); assert.equal(unavailable.code, "search_unavailable"); const defect = await runWeb(Effect.die(new Error(secret))); assert.equal(defect.code, "request_failed"); assert.doesNotMatch( JSON.stringify([unavailable, defect]), new RegExp(secret), ); }); test( "native shell SSE consumption keeps the scope open and closes on every exit", { timeout: 20_000 }, async () => { const { Miniflare, convertV4MiniflareOptions } = await import("miniflare"); const native = await build({ entryPoints: ["tests/fixtures/effect-shell-worker.ts"], bundle: true, write: false, format: "esm", platform: "neutral", conditions: ["workerd", "worker", "browser"], mainFields: ["module", "main"], external: ["cloudflare:*", "node:*"], }); const worker = new Miniflare( convertV4MiniflareOptions({ modules: true, script: native.outputFiles[0].text, compatibilityDate: "2026-09-04", compatibilityFlags: ["nodejs_compat"], }), ); try { const run = async (scenario) => (await worker.dispatchFetch(`http://fixture/${scenario}`)).json(); const success = await run("success"); assert.equal(success.results.at(-1).status, "succeeded"); assert.equal(success.results.at(-1).stdout, "Hello"); assert.equal(success.results.at(-1).cleanup, "closed"); assert.equal(success.closed, 1); const early = await run("early-return"); assert.equal(early.results[0].status, "running"); assert.equal(early.closed, 1); assert.equal(early.launched, 0); assert.deepEqual(await run("return-during-reserve"), { reserved: 1, launched: 0, closed: 1, cancelled: 0, }); const pre = await run("pre-abort"); assert.equal(pre.reserved, 0); assert.equal(pre.results.at(-1).error, "cancelled"); const late = await run("late-reserve"); assert.equal(late.closed, 1); assert.equal(late.launched, 0); assert.equal(late.results.at(-1).error, "cancelled"); const cancel = await run("cancel-body"); assert.equal(cancel.closed, 1); assert.equal(cancel.cancelled, 1); assert.equal(cancel.results.at(-1).error, "cancelled"); const limit = await run("output-limit"); assert.equal(limit.closed, 1); assert.equal(limit.cancelled, 1); assert.equal(limit.results.at(-1).error, "output_limit"); assert.equal(limit.results.at(-1).outputBytes, 32768); assert.equal(limit.results.at(-1).stdout.includes("�"), false); } finally { await worker.dispose(); } }, ); test( "an unresponsive body cancellation cannot hold the request scope open", { timeout: 2000 }, async () => { const reading = Promise.withResolvers(); let cancelled = 0; const body = new ReadableStream({ pull() { reading.resolve(); }, cancel() { cancelled++; return new Promise(() => {}); }, }); const caller = new AbortController(); const running = runWeb( searchWeb({ run: async () => new Response(body) }, "query", 3), caller.signal, ); await reading.promise; caller.abort(); assert.equal((await running).code, "cancelled"); assert.equal(cancelled, 1); assert.equal(body.locked, false); }, ); test( "native bootstrap health closes the SSE scope and caches only confirmed teardown", { timeout: 20_000 }, async () => { const { Miniflare, convertV4MiniflareOptions } = await import("miniflare"); const native = await build({ entryPoints: ["tests/fixtures/effect-health-worker.ts"], bundle: true, write: false, format: "esm", platform: "neutral", conditions: ["workerd", "worker", "browser"], mainFields: ["module", "main"], external: ["cloudflare:*", "node:*"], }); const worker = new Miniflare( convertV4MiniflareOptions({ modules: true, script: native.outputFiles[0].text, compatibilityDate: "2026-09-04", compatibilityFlags: ["nodejs_compat"], durableObjects: { HEALTH: { className: "HealthFixture", useSQLite: true }, }, }), ); try { const run = async (scenario) => (await worker.dispatchFetch(`http://fixture/${scenario}`)).json(); assert.deepEqual(await run("success"), { ok: true, closed: 1, cancelled: 1, pending: 0, }); assert.deepEqual(await run("invalid-output"), { ok: false, closed: 1, cancelled: 1, pending: 0, }); assert.deepEqual(await run("pending-close"), { ok: false, closed: 1, cancelled: 1, pending: 1, }); } finally { await worker.dispose(); } }, ); test("native deadlines, late browser cleanup and declared RPC failures retain their boundaries", async () => { const { Miniflare, convertV4MiniflareOptions } = await import("miniflare"); const { customerBindings } = await import("./fixtures/config.mjs"); const native = await build({ entryPoints: ["tests/fixtures/effect-boundaries-worker.ts"], alias: { path: "node:path" }, bundle: true, write: false, format: "esm", platform: "neutral", conditions: ["workerd", "worker", "browser"], mainFields: ["module", "main"], external: ["cloudflare:*", "node:*"], }); const worker = new Miniflare( convertV4MiniflareOptions({ modules: true, script: native.outputFiles[0].text, compatibilityDate: "2026-09-04", compatibilityFlags: ["nodejs_compat"], bindings: customerBindings, durableObjects: { FAILURES: { className: "FailureFixture", useSQLite: true }, }, }), ); try { const deadlines = await ( await worker.dispatchFetch("http://fixture/deadlines") ).json(); assert.deepEqual(deadlines, { success: 0, rejected: 0, thrown: 0, interrupted: 0, }); for (const scenario of [ "success", "pre-abort", "late-create", "late-connect", "cancel-command", "timeout-command", "late-create-register-fail", "register-delay", "register-fail", "metadata-stall", ]) { const race = await ( await worker.dispatchFetch( `http://fixture/browser-race?scenario=${scenario}`, ) ).json(); assert.equal( race.deletes.length, scenario === "pre-abort" ? 0 : 1, scenario, ); assert.equal( race.result, scenario === "success" || scenario === "metadata-stall", scenario, ); if (scenario === "pre-abort") assert.equal(race.creates, 0); } const result = await ( await worker.dispatchFetch("http://fixture/rpc") ).json(); assert.deepEqual(result.raw, { ok: false, error: "credential_unreadable" }); assert.deepEqual(result.failures, [ "Stored provider API key cannot be read; replace it in settings", "Memory changed or was deleted. Reload memories before editing.", "Personal runtime unavailable", ]); assert.equal(result.model.configuration.provider, "workers-ai"); assert.equal(result.memory.content, "Prefers tea"); assert.equal(result.memory.version, 2); assert.ok(!JSON.stringify(result).includes("PRIVATE-RPC-FAILURE")); } finally { await worker.dispose(); } });