Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370import 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(); }});