From f8bb727b79fcb3edc12e2e44539cc0b6016e9d92 Mon Sep 17 00:00:00 2001 From: Rhys Sullivan <39114868+RhysSullivan@users.noreply.github.com> Date: Wed, 7 Oct 2026 21:35:11 -0700 Subject: [PATCH 1/6] Test Durable Object failure reports and the reset retry --- .../cloudflare/src/__tests__/worker.test.ts | 115 ++++++++++++++++++ 1 file changed, 115 insertions(+) diff --git a/packages/@emulators/cloudflare/src/__tests__/worker.test.ts b/packages/@emulators/cloudflare/src/__tests__/worker.test.ts index 00874153..f6962132 100644 --- a/packages/@emulators/cloudflare/src/__tests__/worker.test.ts +++ b/packages/@emulators/cloudflare/src/__tests__/worker.test.ts @@ -252,6 +252,97 @@ describe("cloudflare worker routing", () => { }); }); +describe("cloudflare worker durable object failures", () => { + // Cloudflare raises Durable Object stub failures as exceptions carrying + // `.retryable`, `.overloaded` and `.remote` flags. + const doError = (message: string, flags: { retryable?: boolean; overloaded?: boolean; remote?: boolean }) => + Object.assign(new Error(message), flags); + + const failingEnv = (failures: Error[]) => { + const calls: string[] = []; + const env: Env = { + EMULATOR: { + idFromName: (n) => n, + get: () => ({ + async fetch(request) { + calls.push(`${request.method} ${new URL(request.url).pathname}`); + const failure = failures.shift(); + if (failure) throw failure; + return Response.json({ ok: true }); + }, + }), + }, + }; + return { env, calls }; + }; + + it("retries a reset that Cloudflare marks retryable, with a fresh stub", async () => { + const { env, calls } = failingEnv([doError("Network connection lost.", { retryable: true })]); + const response = await worker.fetch( + new Request("https://emulators.dev/resend/run-1/_emulate/reset", { method: "POST", body: "{}" }), + env, + ); + expect(response.status).toBe(200); + expect(calls).toEqual(["POST /_emulate/reset", "POST /_emulate/reset"]); + }); + + it("does not retry a credential request; it names the endpoint, instance and cause", async () => { + const { env, calls } = failingEnv([doError("Network connection lost.", { retryable: true })]); + const response = await worker.fetch( + new Request("https://emulators.dev/google/run-1/_emulate/credentials", { + method: "POST", + headers: { "cf-ray": "abc123-PHX" }, + body: JSON.stringify({ type: "api-key" }), + }), + env, + ); + expect(calls).toEqual(["POST /_emulate/credentials"]); + expect(response.status).toBe(503); + expect(await response.json()).toEqual({ + error: "emulator_unavailable", + message: "Network connection lost.", + service: "google", + instance: "run-1", + method: "POST", + path: "/_emulate/credentials", + attempts: 1, + retryable: true, + overloaded: false, + remote: false, + ray: "abc123-PHX", + }); + }); + + it("never retries an overloaded object", async () => { + const { env, calls } = failingEnv([ + doError("Durable Object is overloaded. Too many requests queued.", { retryable: true, overloaded: true }), + ]); + const response = await worker.fetch(new Request("https://emulators.dev/resend/run-1/emails"), env); + expect(calls).toEqual(["GET /emails"]); + expect(response.status).toBe(503); + expect(await response.json()).toMatchObject({ overloaded: true, attempts: 1 }); + }); + + it("stops after a bounded number of attempts", async () => { + const lost = () => doError("Network connection lost.", { retryable: true }); + const { env, calls } = failingEnv([lost(), lost(), lost(), lost()]); + const response = await worker.fetch(new Request("https://emulators.dev/resend/run-1/emails"), env); + expect(calls).toHaveLength(3); + expect(response.status).toBe(503); + expect(await response.json()).toMatchObject({ attempts: 3, path: "/emails" }); + }); + + it("reports an object killed by its own limits as a 500 without retrying", async () => { + const { env, calls } = failingEnv([ + doError("Durable Object's isolate exceeded its memory limit and was reset.", { remote: true }), + ]); + const response = await worker.fetch(new Request("https://emulators.dev/resend/run-1/emails"), env); + expect(calls).toHaveLength(1); + expect(response.status).toBe(500); + expect(await response.json()).toMatchObject({ remote: true, retryable: false, service: "resend" }); + }); +}); + describe("cloudflare durable object control plane", () => { function makeState(options: { limit?: number; initial?: Record } = {}) { const storage = new Map(); @@ -351,6 +442,30 @@ describe("cloudflare durable object control plane", () => { ...extra, }); + it("answers an emulator failure with its cause instead of throwing", async () => { + // A 1-byte value cap makes every persist fail the way an oversized value does. + const { state } = makeState({ limit: 1 }); + const durableObject = new EmulatorDurableObject(state, {}); + const response = await durableObject.fetch( + new Request("https://github.my-run.emulators.dev/_emulate/reset", { + method: "POST", + headers: { ...idHeaders({ "cf-ray": "abc123-SJC" }), "content-type": "application/json" }, + body: "{}", + }), + ); + expect(response.status).toBe(500); + const body = (await response.json()) as Record; + expect(body).toMatchObject({ + error: "emulator_error", + service: "github", + instance: "my-run", + method: "POST", + path: "/_emulate/reset", + ray: "abc123-SJC", + }); + expect(String(body.message)).toMatch(/^Values cannot be larger than 1 bytes/); + }); + // Executor's cloud onboarding e2e provisions this service exactly this way: // mint an api-key, seed a brand, then resolve the company from a work email. // A missing registration only shows up here, as a 404 from the control plane. From b82a4fc15621e799139d64b87b6217baaf0369b5 Mon Sep 17 00:00:00 2001 From: Rhys Sullivan <39114868+RhysSullivan@users.noreply.github.com> Date: Wed, 7 Oct 2026 21:52:17 -0700 Subject: [PATCH 2/6] Test that Cloudflare's retryable storage failures pass through --- .../cloudflare/src/__tests__/worker.test.ts | 20 +++++++++++++++++++ 1 file changed, 20 insertions(+) diff --git a/packages/@emulators/cloudflare/src/__tests__/worker.test.ts b/packages/@emulators/cloudflare/src/__tests__/worker.test.ts index f6962132..bc3ab248 100644 --- a/packages/@emulators/cloudflare/src/__tests__/worker.test.ts +++ b/packages/@emulators/cloudflare/src/__tests__/worker.test.ts @@ -466,6 +466,26 @@ describe("cloudflare durable object control plane", () => { expect(String(body.message)).toMatch(/^Values cannot be larger than 1 bytes/); }); + it("lets Cloudflare's retryable storage failures reach the Worker with their flags", async () => { + const { state } = makeState(); + const moved = Object.assign(new Error("cannot access storage because object has moved to a different machine"), { + retryable: true, + }); + state.storage.get = async () => { + throw moved; + }; + const durableObject = new EmulatorDurableObject(state, {}); + await expect( + durableObject.fetch( + new Request("https://github.my-run.emulators.dev/_emulate/reset", { + method: "POST", + headers: { ...idHeaders(), "content-type": "application/json" }, + body: "{}", + }), + ), + ).rejects.toBe(moved); + }); + // Executor's cloud onboarding e2e provisions this service exactly this way: // mint an api-key, seed a brand, then resolve the company from a work email. // A missing registration only shows up here, as a 404 from the control plane. From 28178fcaebc4cb94398da3a0f98d932a94379f52 Mon Sep 17 00:00:00 2001 From: Rhys Sullivan <39114868+RhysSullivan@users.noreply.github.com> Date: Wed, 7 Oct 2026 22:45:16 -0700 Subject: [PATCH 3/6] Test redacted failure reports and flags carried through the router Replaces the retry tests: a retryable reset and a WorkOS authorize are answered once. Synthetic tokens, codes, addresses and the instance name in an error, path or header must not reach the response or the log. A one-shot storage failure inside the router, the control plane's catches and the object's load keeps its flags all the way to the Worker's 503. --- .../cloudflare/src/__tests__/worker.test.ts | 301 ++++++++++++++---- .../core/src/__tests__/control-plane.test.ts | 30 ++ .../core/src/__tests__/http.test.ts | 28 ++ 3 files changed, 294 insertions(+), 65 deletions(-) diff --git a/packages/@emulators/cloudflare/src/__tests__/worker.test.ts b/packages/@emulators/cloudflare/src/__tests__/worker.test.ts index bc3ab248..206d2ac9 100644 --- a/packages/@emulators/cloudflare/src/__tests__/worker.test.ts +++ b/packages/@emulators/cloudflare/src/__tests__/worker.test.ts @@ -1,5 +1,6 @@ import { describe, expect, it, vi } from "vitest"; import { EmulatorDurableObject } from "../durable-object.js"; +import { instanceId } from "../diagnostics.js"; import worker, { parseHostRoute, type Env } from "../worker.js"; describe("cloudflare worker routing", () => { @@ -252,6 +253,32 @@ describe("cloudflare worker routing", () => { }); }); +// Synthetic secrets: an instance URL is the only access control for its +// emulator, and errors and paths can quote tokens, codes and addresses. None of +// these may reach a failure response or a log line. +const SECRET_INSTANCE = "d040-probe-0123456789abcdef01234567"; +const SECRETS = [ + SECRET_INSTANCE, + "0123456789abcdef01234567", + "emu_resend_SYNTHETICtoken0001", + "SYNTH-CODE-482913", + "pat.synthetic@example.test", +]; +const leakedSecrets = (text: string) => SECRETS.filter((secret) => text.includes(secret)); +const REPORT_KEYS = [ + "error", + "errorClass", + "instanceId", + "method", + "overloaded", + "ray", + "remote", + "retryable", + "route", + "service", +]; +const RAY = "8f1d2c3b4a5e6f70-PHX"; + describe("cloudflare worker durable object failures", () => { // Cloudflare raises Durable Object stub failures as exceptions carrying // `.retryable`, `.overloaded` and `.remote` flags. @@ -276,67 +303,124 @@ describe("cloudflare worker durable object failures", () => { return { env, calls }; }; - it("retries a reset that Cloudflare marks retryable, with a fresh stub", async () => { - const { env, calls } = failingEnv([doError("Network connection lost.", { retryable: true })]); - const response = await worker.fetch( - new Request("https://emulators.dev/resend/run-1/_emulate/reset", { method: "POST", body: "{}" }), - env, - ); - expect(response.status).toBe(200); - expect(calls).toEqual(["POST /_emulate/reset", "POST /_emulate/reset"]); + const captureErrors = () => { + const spy = vi.spyOn(console, "error").mockImplementation(() => {}); + return { + text: () => spy.mock.calls.map((args) => args.map(String).join(" ")).join("\n"), + restore: () => spy.mockRestore(), + }; + }; + + it("answers a retryable reset once, with a report and no replay", async () => { + const logs = captureErrors(); + try { + const { env, calls } = failingEnv([doError("Network connection lost.", { retryable: true })]); + const response = await worker.fetch( + new Request(`https://emulators.dev/resend/${SECRET_INSTANCE}/_emulate/reset`, { + method: "POST", + headers: { "cf-ray": RAY }, + body: "{}", + }), + env, + ); + // A retryable flag does not prove the object never ran the reset. + expect(calls).toEqual(["POST /_emulate/reset"]); + expect(response.status).toBe(503); + const report = await response.json(); + expect(report).toEqual({ + error: "emulator_unavailable", + service: "resend", + instanceId: await instanceId("resend", SECRET_INSTANCE), + method: "POST", + route: "/_emulate/reset", + errorClass: "Error", + retryable: true, + overloaded: false, + remote: false, + ray: RAY, + }); + expect(JSON.parse(logs.text())).toEqual(report); + } finally { + logs.restore(); + } }); - it("does not retry a credential request; it names the endpoint, instance and cause", async () => { + it("does not replay a WorkOS authorize, which issues a code on every call", async () => { const { env, calls } = failingEnv([doError("Network connection lost.", { retryable: true })]); const response = await worker.fetch( - new Request("https://emulators.dev/google/run-1/_emulate/credentials", { - method: "POST", - headers: { "cf-ray": "abc123-PHX" }, - body: JSON.stringify({ type: "api-key" }), - }), + new Request( + `https://emulators.dev/workos/${SECRET_INSTANCE}/oauth2/authorize?client_id=c&redirect_uri=https%3A%2F%2Fapp.example.test%2Fcb`, + ), env, ); - expect(calls).toEqual(["POST /_emulate/credentials"]); + expect(calls).toEqual(["GET /oauth2/authorize"]); expect(response.status).toBe(503); - expect(await response.json()).toEqual({ - error: "emulator_unavailable", - message: "Network connection lost.", - service: "google", - instance: "run-1", - method: "POST", - path: "/_emulate/credentials", - attempts: 1, - retryable: true, - overloaded: false, - remote: false, - ray: "abc123-PHX", - }); + expect(await response.json()).toMatchObject({ route: "/oauth2/authorize", retryable: true }); + }); + + it("keeps tokens, codes, addresses and the instance out of the response and the log", async () => { + const logs = captureErrors(); + try { + const failure = doError( + `lost while serving ${SECRET_INSTANCE}: token emu_resend_SYNTHETICtoken0001 code SYNTH-CODE-482913 for pat.synthetic@example.test`, + { retryable: true }, + ); + const { env } = failingEnv([failure]); + const response = await worker.fetch( + new Request( + `https://emulators.dev/resend/${SECRET_INSTANCE}/domains/pat.synthetic@example.test?code=SYNTH-CODE-482913`, + { + headers: { authorization: "Bearer emu_resend_SYNTHETICtoken0001", "cf-ray": RAY }, + }, + ), + env, + ); + const text = await response.text(); + const report = JSON.parse(text) as Record; + expect(Object.keys(report).sort()).toEqual(REPORT_KEYS); + expect(report.route).toBe("/domains/:id"); + expect(leakedSecrets(text)).toEqual([]); + expect(logs.text()).not.toBe(""); + expect(leakedSecrets(logs.text())).toEqual([]); + } finally { + logs.restore(); + } + }); + + it("does not echo header or path values it cannot vouch for", async () => { + const logs = captureErrors(); + try { + const { env } = failingEnv([doError("Network connection lost.", { retryable: true })]); + const response = await worker.fetch( + new Request(`https://emulators.dev/pat.synthetic@example.test/${SECRET_INSTANCE}/emails`, { + headers: { "cf-ray": "SYNTH-CODE-482913" }, + }), + env, + ); + const text = await response.text(); + expect(JSON.parse(text)).toMatchObject({ service: "unknown", route: "unmatched", ray: null }); + expect(leakedSecrets(text)).toEqual([]); + expect(leakedSecrets(logs.text())).toEqual([]); + } finally { + logs.restore(); + } }); - it("never retries an overloaded object", async () => { + it("answers an overloaded object with a 503", async () => { const { env, calls } = failingEnv([ doError("Durable Object is overloaded. Too many requests queued.", { retryable: true, overloaded: true }), ]); - const response = await worker.fetch(new Request("https://emulators.dev/resend/run-1/emails"), env); + const response = await worker.fetch(new Request(`https://emulators.dev/resend/${SECRET_INSTANCE}/emails`), env); expect(calls).toEqual(["GET /emails"]); expect(response.status).toBe(503); - expect(await response.json()).toMatchObject({ overloaded: true, attempts: 1 }); + expect(await response.json()).toMatchObject({ overloaded: true, route: "/emails" }); }); - it("stops after a bounded number of attempts", async () => { - const lost = () => doError("Network connection lost.", { retryable: true }); - const { env, calls } = failingEnv([lost(), lost(), lost(), lost()]); - const response = await worker.fetch(new Request("https://emulators.dev/resend/run-1/emails"), env); - expect(calls).toHaveLength(3); - expect(response.status).toBe(503); - expect(await response.json()).toMatchObject({ attempts: 3, path: "/emails" }); - }); - - it("reports an object killed by its own limits as a 500 without retrying", async () => { + it("reports an object killed by its own limits as a 500", async () => { const { env, calls } = failingEnv([ doError("Durable Object's isolate exceeded its memory limit and was reset.", { remote: true }), ]); - const response = await worker.fetch(new Request("https://emulators.dev/resend/run-1/emails"), env); + const response = await worker.fetch(new Request(`https://emulators.dev/resend/${SECRET_INSTANCE}/emails`), env); expect(calls).toHaveLength(1); expect(response.status).toBe(500); expect(await response.json()).toMatchObject({ remote: true, retryable: false, service: "resend" }); @@ -442,48 +526,135 @@ describe("cloudflare durable object control plane", () => { ...extra, }); - it("answers an emulator failure with its cause instead of throwing", async () => { + // Builds the instance, then makes the next storage write throw `failure` once. + // One-shot, so a later write (the persist after every mutating request) + // cannot rethrow it and hide what the router did with the first one. + const failNextWriteAfterWarmup = async (failure: unknown) => { + const { state } = makeState(); + const durableObject = new EmulatorDurableObject(state, {}); + const warm = await durableObject.fetch( + new Request("https://github.my-run.emulators.dev/_emulate/manifest", { headers: idHeaders() }), + ); + expect(warm.status).toBe(200); + const put = state.storage.put; + let thrown = false; + state.storage.put = async (key, value) => { + if (thrown) return put(key, value); + thrown = true; + throw failure; + }; + return durableObject; + }; + const control = (path: string, body: unknown, extra: Record = {}) => + new Request(`https://github.my-run.emulators.dev${path}`, { + method: "POST", + headers: { ...idHeaders(extra), "content-type": "application/json" }, + body: JSON.stringify(body), + }); + const moved = () => + Object.assign(new Error("cannot access storage because object has moved to a different machine"), { + retryable: true, + }); + + it("answers an emulator failure with a report instead of throwing", async () => { // A 1-byte value cap makes every persist fail the way an oversized value does. const { state } = makeState({ limit: 1 }); const durableObject = new EmulatorDurableObject(state, {}); - const response = await durableObject.fetch( - new Request("https://github.my-run.emulators.dev/_emulate/reset", { - method: "POST", - headers: { ...idHeaders({ "cf-ray": "abc123-SJC" }), "content-type": "application/json" }, - body: "{}", - }), - ); + const response = await durableObject.fetch(control("/_emulate/reset", {}, { "cf-ray": RAY })); expect(response.status).toBe(500); - const body = (await response.json()) as Record; - expect(body).toMatchObject({ + expect(await response.json()).toEqual({ error: "emulator_error", service: "github", - instance: "my-run", + instanceId: await instanceId("github", "my-run"), method: "POST", - path: "/_emulate/reset", - ray: "abc123-SJC", + route: "/_emulate/reset", + errorClass: "Error", + retryable: false, + overloaded: false, + remote: false, + ray: RAY, }); - expect(String(body.message)).toMatch(/^Values cannot be larger than 1 bytes/); + }); + + it("keeps secrets in an emulator error out of the response and the log", async () => { + const logs = vi.spyOn(console, "error").mockImplementation(() => {}); + try { + const durableObject = await failNextWriteAfterWarmup( + new Error( + `write failed for ${SECRET_INSTANCE}: token emu_resend_SYNTHETICtoken0001 code SYNTH-CODE-482913 for pat.synthetic@example.test`, + ), + ); + // The error is thrown inside the router (reset runs in the control plane), + // which used to answer it with its raw message. + const response = await durableObject.fetch(control("/_emulate/reset", {})); + expect(response.status).toBe(500); + const text = await response.text(); + expect(Object.keys(JSON.parse(text)).sort()).toEqual(REPORT_KEYS); + expect(JSON.parse(text)).toMatchObject({ error: "emulator_error", route: "/_emulate/reset" }); + expect(leakedSecrets(text)).toEqual([]); + const logged = logs.mock.calls.map((args) => args.map(String).join(" ")).join("\n"); + expect(logged).not.toBe(""); + expect(leakedSecrets(logged)).toEqual([]); + } finally { + logs.mockRestore(); + } }); it("lets Cloudflare's retryable storage failures reach the Worker with their flags", async () => { const { state } = makeState(); - const moved = Object.assign(new Error("cannot access storage because object has moved to a different machine"), { - retryable: true, - }); + const failure = moved(); state.storage.get = async () => { - throw moved; + throw failure; }; const durableObject = new EmulatorDurableObject(state, {}); + // Thrown while the object loads, before the service router runs. + await expect(durableObject.fetch(control("/_emulate/reset", {}))).rejects.toBe(failure); + }); + + it("carries a flagged failure through the service router", async () => { + const failure = moved(); + const durableObject = await failNextWriteAfterWarmup(failure); + // Reset persists from inside the router, whose error handler used to answer + // every error as a plain 500 without the flags. + await expect(durableObject.fetch(control("/_emulate/reset", {}))).rejects.toBe(failure); + }); + + it("carries a flagged failure through the control plane's own catches", async () => { + const credentials = moved(); await expect( - durableObject.fetch( - new Request("https://github.my-run.emulators.dev/_emulate/reset", { + (await failNextWriteAfterWarmup(credentials)).fetch( + control("/_emulate/credentials", { type: "bearer-token", login: "synthetic-user" }), + ), + ).rejects.toBe(credentials); + const seed = moved(); + await expect( + (await failNextWriteAfterWarmup(seed)).fetch(control("/_emulate/seed", { users: [{ login: "synthetic-user" }] })), + ).rejects.toBe(seed); + }); + + it("returns the flagged failure to the Worker, which reports it as a 503", async () => { + const logs = vi.spyOn(console, "error").mockImplementation(() => {}); + try { + const durableObject = await failNextWriteAfterWarmup(moved()); + const env: Env = { EMULATOR: { idFromName: (n) => n, get: () => durableObject } }; + const response = await worker.fetch( + new Request("https://emulators.dev/github/my-run/_emulate/reset", { method: "POST", - headers: { ...idHeaders(), "content-type": "application/json" }, + headers: { "cf-ray": RAY }, body: "{}", }), - ), - ).rejects.toBe(moved); + env, + ); + expect(response.status).toBe(503); + expect(await response.json()).toMatchObject({ + error: "emulator_unavailable", + route: "/_emulate/reset", + retryable: true, + ray: RAY, + }); + } finally { + logs.mockRestore(); + } }); // Executor's cloud onboarding e2e provisions this service exactly this way: diff --git a/packages/@emulators/core/src/__tests__/control-plane.test.ts b/packages/@emulators/core/src/__tests__/control-plane.test.ts index a2649697..73006fe7 100644 --- a/packages/@emulators/core/src/__tests__/control-plane.test.ts +++ b/packages/@emulators/core/src/__tests__/control-plane.test.ts @@ -3,6 +3,7 @@ import { createServer } from "../server.js"; import { randomInstanceName } from "../control-plane.js"; import type { ServicePlugin } from "../plugin.js"; import type { Store } from "../store.js"; +import { ApiError } from "../middleware/error-handler.js"; interface Thing { id: number; @@ -307,3 +308,32 @@ describe("randomInstanceName", () => { expect(name).toMatch(/-[0-9a-f]{24}$/); }); }); + +describe("unexpected route errors", () => { + const throwing: ServicePlugin = { + name: "throwing", + register(app) { + app.get("/missing", () => { + throw new ApiError(404, "Thing not found"); + }); + app.get("/broken", () => { + throw new Error("token emu_demo_SYNTHETIC0001 for pat.synthetic@example.test"); + }); + }, + }; + + it("answers them with their message by default", async () => { + const { app } = createServer(throwing); + const res = await app.request("/broken"); + expect(res.status).toBe(500); + expect(await res.json()).toMatchObject({ message: "token emu_demo_SYNTHETIC0001 for pat.synthetic@example.test" }); + }); + + it("rethrows them for a host that reports failures itself, but keeps API errors", async () => { + const { app } = createServer(throwing, { rethrowUnexpectedErrors: true }); + await expect(app.request("/broken")).rejects.toThrow("emu_demo_SYNTHETIC0001"); + const res = await app.request("/missing"); + expect(res.status).toBe(404); + expect(await res.json()).toMatchObject({ message: "Thing not found" }); + }); +}); diff --git a/packages/@emulators/core/src/__tests__/http.test.ts b/packages/@emulators/core/src/__tests__/http.test.ts index 82421938..0d73c039 100644 --- a/packages/@emulators/core/src/__tests__/http.test.ts +++ b/packages/@emulators/core/src/__tests__/http.test.ts @@ -71,4 +71,32 @@ describe("internal http layer", () => { expect(res.headers.get("Access-Control-Allow-Headers")).toBe("x-test"); expect(res.headers.get("Access-Control-Max-Age")).toBe("60"); }); + + it("passes Cloudflare's flagged failures through the error handler", async () => { + const moved = Object.assign(new Error("object has moved to a different machine"), { retryable: true }); + const overloaded = Object.assign(new Error("Durable Object is overloaded."), { overloaded: true }); + const app = new Hono(); + app.onError(() => new Response("handled", { status: 500 })); + app.get("/moved", () => { + throw moved; + }); + app.get("/overloaded", () => { + throw overloaded; + }); + app.get("/plain", () => { + throw new Error("plain"); + }); + + await expect(app.request("/moved")).rejects.toBe(moved); + await expect(app.request("/overloaded")).rejects.toBe(overloaded); + expect(await (await app.request("/plain")).text()).toBe("handled"); + }); + + it("names the route pattern a request would match", () => { + const app = new Hono(); + app.get("/domains/:id", (c) => c.text("ok")); + expect(app.routePattern("GET", "/domains/pat.synthetic@example.test")).toBe("/domains/:id"); + expect(app.routePattern("HEAD", "/domains/x")).toBe("/domains/:id"); + expect(app.routePattern("POST", "/domains/x")).toBeUndefined(); + }); }); From f93ca4e27c375462e66f3ae3a7250bb1feefe9d0 Mon Sep 17 00:00:00 2001 From: Rhys Sullivan <39114868+RhysSullivan@users.noreply.github.com> Date: Thu, 8 Oct 2026 00:31:56 -0700 Subject: [PATCH 4/6] Test that no Worker or Durable Object failure escapes, in workerd too A Miniflare probe runs the real Worker and Durable Object with injected storage, addressing and stub failures, and checks every tail event, workerd's output and each response: no exception events and no secret. Unit tests cover secret-bearing error names and methods, body and addressing failures, and plain seed and credential failures. --- packages/@emulators/cloudflare/package.json | 2 + .../cloudflare/src/__tests__/runtime.test.ts | 220 ++++++++++++++++++ .../cloudflare/src/__tests__/worker.test.ts | 201 ++++++++++++++-- .../core/src/__tests__/control-plane.test.ts | 40 ++++ .../core/src/__tests__/http.test.ts | 20 -- pnpm-lock.yaml | 6 + 6 files changed, 449 insertions(+), 40 deletions(-) create mode 100644 packages/@emulators/cloudflare/src/__tests__/runtime.test.ts diff --git a/packages/@emulators/cloudflare/package.json b/packages/@emulators/cloudflare/package.json index a5ea70ec..7488d0c4 100644 --- a/packages/@emulators/cloudflare/package.json +++ b/packages/@emulators/cloudflare/package.json @@ -52,6 +52,8 @@ "@emulators/planetscale": "workspace:*" }, "devDependencies": { + "esbuild": "0.27.4", + "miniflare": "4.20260702.0", "tsup": "^8", "typescript": "^5.7" } diff --git a/packages/@emulators/cloudflare/src/__tests__/runtime.test.ts b/packages/@emulators/cloudflare/src/__tests__/runtime.test.ts new file mode 100644 index 00000000..e7a49a67 --- /dev/null +++ b/packages/@emulators/cloudflare/src/__tests__/runtime.test.ts @@ -0,0 +1,220 @@ +import { fileURLToPath } from "node:url"; +import { isBuiltin } from "node:module"; +import type { Readable } from "node:stream"; +import { build } from "esbuild"; +import { Log, LogLevel, Miniflare } from "miniflare"; +import { afterAll, beforeAll, describe, expect, it } from "vitest"; + +// Runs the real Worker and Durable Object in workerd and records everything +// Cloudflare would: every tail event (the source of Workers Logs, including +// uncaught exception events) and workerd's own stdout and stderr. Faults are +// injected at the storage and stub boundaries with synthetic secrets in the +// message, the error name and the instance URL; none may reach any of it. +const SECRET = "SYNTHETICtoken482913"; +const SUFFIX = "0123456789abcdef01234567"; +const INSTANCE = `synthetic-${SUFFIX}`; +const SECRETS = [SECRET, SUFFIX]; + +// The probe module wraps the shipped exports. Its Durable Object hands the real +// one a storage proxy that fails the way `x-probe-fault` asks, and its Worker +// can swap in a namespace whose addressing or stub fails. +const PROBE = ` +import worker, { EmulatorDurableObject } from "./worker.ts"; +const secretError = (flags) => + Object.assign(new Error("uncaught ${SECRET} https://resend.${INSTANCE}.emulators.dev/emails"), { name: "${SECRET}" }, flags); +export class ProbeObject extends EmulatorDurableObject { + constructor(state, env) { + let fault = null; + const storage = new Proxy(state.storage, { + get(target, key) { + if (fault === "do-flagged" && key === "get") return async () => { throw secretError({ retryable: true }); }; + if (fault === "do-plain" && key === "get") return async () => { throw secretError({}); }; + if (fault === "do-put" && key === "put") return async () => { throw secretError({}); }; + const value = target[key]; + return typeof value === "function" ? value.bind(target) : value; + }, + }); + super({ storage, blockConcurrencyWhile: state.blockConcurrencyWhile.bind(state) }, env); + this.setFault = (next) => { fault = next; }; + } + async fetch(request) { + this.setFault(request.headers.get("x-probe-fault")); + return super.fetch(request); + } +} +export default { + fetch(request, env) { + const fault = request.headers.get("x-probe-fault"); + if (fault === "worker-id") + env = { ...env, EMULATOR: { idFromName() { throw secretError({}); }, get: () => env.EMULATOR.get() } }; + if (fault === "worker-stub") + env = { ...env, EMULATOR: { idFromName: (n) => n, get: () => ({ fetch: async () => { throw secretError({ retryable: true, overloaded: true }); } }) } }; + if (fault === "worker-env") + env = { get EMULATE_HOST_SUFFIX() { throw secretError({}); } }; + return worker.fetch(request, env); + }, +}; +`; + +interface TailEvent { + outcome: string; + event?: unknown; + entrypoint?: string; + exceptions: Array<{ name: string; message: string; stack?: string }>; + logs: Array<{ level: string; message: unknown[] }>; +} + +let mf: Miniflare; +const tailed: TailEvent[] = []; +const runtimeOutput: string[] = []; + +beforeAll(async () => { + const bundle = await build({ + stdin: { contents: PROBE, resolveDir: fileURLToPath(new URL("..", import.meta.url)), loader: "ts" }, + bundle: true, + write: false, + format: "esm", + platform: "neutral", + mainFields: ["module", "main"], + conditions: ["workerd", "worker", "import"], + logLevel: "silent", + plugins: [ + { + name: "node-builtins", + setup(b) { + b.onResolve({ filter: /.*/ }, (args) => + isBuiltin(args.path) ? { path: `node:${args.path.replace(/^node:/, "")}`, external: true } : undefined, + ); + }, + }, + ], + }); + mf = new Miniflare({ + log: new Log(LogLevel.NONE), + handleRuntimeStdio(stdout: Readable, stderr: Readable) { + stdout.on("data", (chunk: Buffer) => runtimeOutput.push(String(chunk))); + stderr.on("data", (chunk: Buffer) => runtimeOutput.push(String(chunk))); + }, + workers: [ + { + name: "emulate-hosts", + modules: [{ type: "ESModule", path: "/probe/worker.mjs", contents: bundle.outputFiles[0].text }], + modulesRoot: "/probe", + compatibilityDate: "2026-06-08", + compatibilityFlags: ["nodejs_compat"], + durableObjects: { EMULATOR: "ProbeObject" }, + bindings: { EMULATE_HOST_SUFFIX: "emulators.dev" }, + tails: ["sink"], + }, + { + name: "sink", + modules: true, + script: `export default { async tail(events, env) { await env.CAPTURE.fetch("https://capture.invalid", { method: "POST", body: JSON.stringify(events) }); } };`, + compatibilityDate: "2026-06-08", + serviceBindings: { + CAPTURE: async (request: Request) => { + tailed.push(...((await request.json()) as TailEvent[])); + return new Response("ok"); + }, + }, + }, + ], + }); + await mf.ready; +}, 60_000); + +afterAll(async () => { + await mf?.dispose(); +}); + +async function waitForTail(count: number): Promise { + const deadline = Date.now() + 5_000; + while (tailed.length < count && Date.now() < deadline) await new Promise((r) => setTimeout(r, 25)); + // Late events (a second invocation, a trailing log) would land here. + await new Promise((r) => setTimeout(r, 200)); +} + +const CASES: Array<{ fault: string; method?: string; path: string; status: number; report: Record }> = + [ + { + fault: "do-flagged", + path: "/emails", + status: 503, + report: { error: "emulator_unavailable", route: "/emails", retryable: true, errorClass: "other" }, + }, + { + fault: "do-plain", + path: "/emails", + status: 500, + report: { error: "emulator_error", route: "/emails", retryable: false, errorClass: "other" }, + }, + { + fault: "do-put", + method: "POST", + path: "/_emulate/seed", + status: 500, + report: { error: "emulator_error", route: "/_emulate/seed", errorClass: "other" }, + }, + { + fault: "do-put", + method: "POST", + path: "/_emulate/credentials", + status: 500, + report: { error: "emulator_error", route: "/_emulate/credentials", errorClass: "other" }, + }, + { + fault: "worker-id", + method: "POST", + path: "/emails", + status: 500, + report: { error: "emulator_unavailable", route: "/emails", errorClass: "other" }, + }, + { + fault: "worker-stub", + path: "/emails", + status: 503, + report: { error: "emulator_unavailable", overloaded: true, errorClass: "other" }, + }, + { + fault: "worker-env", + path: "/emails", + status: 500, + report: { error: "worker_error", service: "unknown", instanceId: null }, + }, + ]; + +describe("emulate-hosts in workerd", () => { + it("records no exception and no secret for any injected failure", async () => { + const responses: string[] = []; + for (const c of CASES) { + const before = tailed.length; + const response = await mf.dispatchFetch(`https://emulators.dev/resend/${INSTANCE}${c.path}`, { + method: c.method ?? "GET", + headers: { "x-probe-fault": c.fault, "content-type": "application/json" }, + body: c.method === "POST" ? JSON.stringify({ type: "api-key", login: "synthetic-user" }) : undefined, + }); + const text = await response.text(); + responses.push(text); + expect({ fault: c.fault, status: response.status }).toEqual({ fault: c.fault, status: c.status }); + expect(JSON.parse(text)).toMatchObject(c.report); + // The Worker's invocation, plus the object's when the request reached it. + await waitForTail(before + 1); + } + + expect(tailed.length).toBeGreaterThanOrEqual(CASES.length); + expect(tailed.map((event) => event.outcome).filter((outcome) => outcome !== "ok")).toEqual([]); + expect(tailed.flatMap((event) => event.exceptions)).toEqual([]); + // Each failure is logged once, as its report. + const logged = tailed.flatMap((event) => event.logs.map((log) => log.message.map(String).join(" "))); + expect(logged).toHaveLength(CASES.length); + // Every event also carries its invocation's request (URL and headers), the + // metadata Workers Logs keeps only as invocation logs, which are off. The + // rest of each event is what the code under test can put there. + const recorded = [ + JSON.stringify(tailed.map(({ event: _event, ...rest }) => rest)), + runtimeOutput.join(""), + ...responses, + ].join("\n"); + expect(SECRETS.filter((secret) => recorded.includes(secret))).toEqual([]); + }, 60_000); +}); diff --git a/packages/@emulators/cloudflare/src/__tests__/worker.test.ts b/packages/@emulators/cloudflare/src/__tests__/worker.test.ts index 206d2ac9..78679ab2 100644 --- a/packages/@emulators/cloudflare/src/__tests__/worker.test.ts +++ b/packages/@emulators/cloudflare/src/__tests__/worker.test.ts @@ -263,6 +263,7 @@ const SECRETS = [ "emu_resend_SYNTHETICtoken0001", "SYNTH-CODE-482913", "pat.synthetic@example.test", + "SYNTHETICtoken482913", ]; const leakedSecrets = (text: string) => SECRETS.filter((secret) => text.includes(secret)); const REPORT_KEYS = [ @@ -425,6 +426,96 @@ describe("cloudflare worker durable object failures", () => { expect(response.status).toBe(500); expect(await response.json()).toMatchObject({ remote: true, retryable: false, service: "resend" }); }); + + it("reports failures to read the body or address the object instead of throwing", async () => { + const logs = captureErrors(); + try { + const secretError = () => new Error(`lost ${SECRET_INSTANCE} token emu_resend_SYNTHETICtoken0001`); + const addressing: Env = { + EMULATOR: { + idFromName: () => { + throw secretError(); + }, + get: () => ({ fetch: async () => Response.json({ ok: true }) }), + }, + }; + const addressed = await worker.fetch( + new Request(`https://emulators.dev/resend/${SECRET_INSTANCE}/emails`), + addressing, + ); + expect(addressed.status).toBe(500); + expect(await addressed.json()).toMatchObject({ error: "emulator_unavailable", route: "/emails" }); + + const { env, calls } = failingEnv([]); + const unreadable = new Request(`https://emulators.dev/resend/${SECRET_INSTANCE}/emails`, { + method: "POST", + body: new ReadableStream({ + start(controller) { + controller.error(secretError()); + }, + }), + duplex: "half", + } as RequestInit); + const read = await worker.fetch(unreadable, env); + expect(calls).toEqual([]); + expect(read.status).toBe(500); + expect(await read.json()).toMatchObject({ error: "emulator_unavailable", route: "/emails" }); + expect(leakedSecrets(logs.text())).toEqual([]); + } finally { + logs.restore(); + } + }); + + it("reports any other Worker failure instead of throwing", async () => { + const logs = captureErrors(); + try { + const env = { + get EMULATE_HOST_SUFFIX(): string { + throw new Error(`config read failed for ${SECRET_INSTANCE}`); + }, + } as unknown as Env; + const response = await worker.fetch(new Request(`https://emulators.dev/resend/${SECRET_INSTANCE}/emails`), env); + const text = await response.text(); + expect(response.status).toBe(500); + expect(JSON.parse(text)).toEqual({ + error: "worker_error", + service: "unknown", + instanceId: null, + method: "GET", + route: "unmatched", + errorClass: "Error", + retryable: false, + overloaded: false, + remote: false, + ray: null, + }); + expect(leakedSecrets(text + logs.text())).toEqual([]); + } finally { + logs.restore(); + } + }); + + it("reports only allowlisted error classes and methods", async () => { + const logs = captureErrors(); + try { + const { env } = failingEnv([ + doError("lost", { retryable: true }), + Object.assign(doError("lost", { retryable: true }), { name: "SYNTHETICtoken482913" }), + Object.assign(new TypeError("lost"), { retryable: true }), + ]); + const send = (method: string) => + worker + .fetch(new Request(`https://emulators.dev/resend/${SECRET_INSTANCE}/emails`, { method }), env) + .then((r) => r.json() as Promise>); + expect(await send("SYNTHETICTOKEN")).toMatchObject({ method: "OTHER", errorClass: "Error" }); + expect(await send("GET")).toMatchObject({ method: "GET", errorClass: "other" }); + expect(await send("GET")).toMatchObject({ errorClass: "TypeError" }); + expect(logs.text()).not.toContain("SYNTHETICTOKEN"); + expect(leakedSecrets(logs.text())).toEqual([]); + } finally { + logs.restore(); + } + }); }); describe("cloudflare durable object control plane", () => { @@ -600,39 +691,107 @@ describe("cloudflare durable object control plane", () => { } }); - it("lets Cloudflare's retryable storage failures reach the Worker with their flags", async () => { + // Every failure below must come back as a report: an error thrown out of the + // object is recorded by Cloudflare with its message, stack and URL. + const reportOf = async (response: Promise) => { + const res = await response; + return { status: res.status, report: (await res.json()) as Record }; + }; + + it("reports Cloudflare's retryable storage failures with their flags instead of throwing", async () => { const { state } = makeState(); - const failure = moved(); state.storage.get = async () => { - throw failure; + throw moved(); }; const durableObject = new EmulatorDurableObject(state, {}); // Thrown while the object loads, before the service router runs. - await expect(durableObject.fetch(control("/_emulate/reset", {}))).rejects.toBe(failure); + expect(await reportOf(durableObject.fetch(control("/_emulate/reset", {}, { "cf-ray": RAY })))).toEqual({ + status: 503, + report: { + error: "emulator_unavailable", + service: "github", + instanceId: await instanceId("github", "my-run"), + method: "POST", + route: "/_emulate/reset", + errorClass: "Error", + retryable: true, + overloaded: false, + remote: false, + ray: RAY, + }, + }); }); - it("carries a flagged failure through the service router", async () => { - const failure = moved(); - const durableObject = await failNextWriteAfterWarmup(failure); + it("reports a flagged failure from inside the service router", async () => { + const durableObject = await failNextWriteAfterWarmup(moved()); // Reset persists from inside the router, whose error handler used to answer // every error as a plain 500 without the flags. - await expect(durableObject.fetch(control("/_emulate/reset", {}))).rejects.toBe(failure); + expect(await reportOf(durableObject.fetch(control("/_emulate/reset", {})))).toMatchObject({ + status: 503, + report: { error: "emulator_unavailable", route: "/_emulate/reset", retryable: true }, + }); }); - it("carries a flagged failure through the control plane's own catches", async () => { - const credentials = moved(); - await expect( - (await failNextWriteAfterWarmup(credentials)).fetch( - control("/_emulate/credentials", { type: "bearer-token", login: "synthetic-user" }), - ), - ).rejects.toBe(credentials); - const seed = moved(); - await expect( - (await failNextWriteAfterWarmup(seed)).fetch(control("/_emulate/seed", { users: [{ login: "synthetic-user" }] })), - ).rejects.toBe(seed); + it("reports flagged and plain storage failures from seed and credential requests", async () => { + const logs = vi.spyOn(console, "error").mockImplementation(() => {}); + try { + const credentials = () => control("/_emulate/credentials", { type: "bearer-token", login: "synthetic-user" }); + const seed = () => control("/_emulate/seed", { users: [{ login: "synthetic-user" }] }); + const plain = () => + new Error( + `storage write failed for ${SECRET_INSTANCE}: token emu_resend_SYNTHETICtoken0001 for pat.synthetic@example.test`, + ); + for (const [request, route] of [ + [credentials, "/_emulate/credentials"], + [seed, "/_emulate/seed"], + ] as const) { + expect(await reportOf((await failNextWriteAfterWarmup(moved())).fetch(request()))).toMatchObject({ + status: 503, + report: { error: "emulator_unavailable", route, retryable: true }, + }); + // A host failure is not the caller's mistake: it used to come back as a + // 400 quoting the raw message. + const res = await (await failNextWriteAfterWarmup(plain())).fetch(request()); + const text = await res.text(); + expect({ status: res.status, report: JSON.parse(text) }).toMatchObject({ + status: 500, + report: { error: "emulator_error", route, errorClass: "Error" }, + }); + expect(leakedSecrets(text)).toEqual([]); + } + expect(leakedSecrets(logs.mock.calls.map((args) => args.map(String).join(" ")).join("\n"))).toEqual([]); + } finally { + logs.mockRestore(); + } + }); + + it("still answers a credential type the emulator does not support with a 400", async () => { + const { state } = makeState(); + const durableObject = new EmulatorDurableObject(state, {}); + const res = await durableObject.fetch(control("/_emulate/credentials", { type: "synthetic-unsupported-type" })); + expect(res.status).toBe(400); + expect(await res.json()).toEqual({ + error: "unsupported", + message: "Credential type synthetic-unsupported-type is not supported by github", + }); + }); + + it("reports an error whose name could carry a secret as class other", async () => { + const logs = vi.spyOn(console, "error").mockImplementation(() => {}); + try { + const named = Object.assign(new Error("synthetic"), { name: "SYNTHETICtoken482913" }); + const durableObject = await failNextWriteAfterWarmup(named); + const res = await durableObject.fetch(control("/_emulate/reset", {})); + const text = await res.text(); + expect(JSON.parse(text)).toMatchObject({ errorClass: "other" }); + expect(text).not.toContain("SYNTHETICtoken482913"); + expect(logs.mock.calls.flat().map(String).join("\n")).not.toContain("SYNTHETICtoken482913"); + } finally { + logs.mockRestore(); + } }); - it("returns the flagged failure to the Worker, which reports it as a 503", async () => { + it("passes the object's flagged report through the Worker as a 503", async () => { const logs = vi.spyOn(console, "error").mockImplementation(() => {}); try { const durableObject = await failNextWriteAfterWarmup(moved()); @@ -652,6 +811,8 @@ describe("cloudflare durable object control plane", () => { retryable: true, ray: RAY, }); + // Reported once, by the object. + expect(logs).toHaveBeenCalledTimes(1); } finally { logs.mockRestore(); } diff --git a/packages/@emulators/core/src/__tests__/control-plane.test.ts b/packages/@emulators/core/src/__tests__/control-plane.test.ts index 73006fe7..bb6915f5 100644 --- a/packages/@emulators/core/src/__tests__/control-plane.test.ts +++ b/packages/@emulators/core/src/__tests__/control-plane.test.ts @@ -4,6 +4,7 @@ import { randomInstanceName } from "../control-plane.js"; import type { ServicePlugin } from "../plugin.js"; import type { Store } from "../store.js"; import { ApiError } from "../middleware/error-handler.js"; +import { ControlPlaneRejection } from "../control-plane-rejection.js"; interface Thing { id: number; @@ -337,3 +338,42 @@ describe("unexpected route errors", () => { expect(await res.json()).toMatchObject({ message: "Thing not found" }); }); }); + +describe("seed and credential failures", () => { + const plugin: ServicePlugin = { name: "seeded", register() {} }; + const failing = (error: unknown, rethrowUnexpectedErrors = false) => + createServer(plugin, { + rethrowUnexpectedErrors, + seed: () => { + throw error; + }, + issueCredential: () => { + throw error; + }, + }).app; + const post = (app: ReturnType, path: string) => + app.request(path, { method: "POST", headers: { "content-type": "application/json" }, body: "{}" }); + + it("answers the emulator's own rejections with a 400 and their message", async () => { + const app = failing(new ControlPlaneRejection("Credential type synthetic is not supported by seeded")); + const seed = await post(app, "/_emulate/seed"); + expect(seed.status).toBe(400); + expect(await seed.json()).toEqual({ + error: "invalid_seed", + message: "Credential type synthetic is not supported by seeded", + }); + const credentials = await post(app, "/_emulate/credentials"); + expect(credentials.status).toBe(400); + expect(await credentials.json()).toMatchObject({ error: "unsupported" }); + }); + + it("sends any other error to the app's error handler, not a 400", async () => { + const error = () => new Error("storage write failed: token emu_demo_SYNTHETIC0001"); + for (const path of ["/_emulate/seed", "/_emulate/credentials"]) { + const res = await post(failing(error()), path); + expect(res.status).toBe(500); + const thrown = error(); + await expect(post(failing(thrown, true), path)).rejects.toBe(thrown); + } + }); +}); diff --git a/packages/@emulators/core/src/__tests__/http.test.ts b/packages/@emulators/core/src/__tests__/http.test.ts index 0d73c039..f616f21f 100644 --- a/packages/@emulators/core/src/__tests__/http.test.ts +++ b/packages/@emulators/core/src/__tests__/http.test.ts @@ -72,26 +72,6 @@ describe("internal http layer", () => { expect(res.headers.get("Access-Control-Max-Age")).toBe("60"); }); - it("passes Cloudflare's flagged failures through the error handler", async () => { - const moved = Object.assign(new Error("object has moved to a different machine"), { retryable: true }); - const overloaded = Object.assign(new Error("Durable Object is overloaded."), { overloaded: true }); - const app = new Hono(); - app.onError(() => new Response("handled", { status: 500 })); - app.get("/moved", () => { - throw moved; - }); - app.get("/overloaded", () => { - throw overloaded; - }); - app.get("/plain", () => { - throw new Error("plain"); - }); - - await expect(app.request("/moved")).rejects.toBe(moved); - await expect(app.request("/overloaded")).rejects.toBe(overloaded); - expect(await (await app.request("/plain")).text()).toBe("handled"); - }); - it("names the route pattern a request would match", () => { const app = new Hono(); app.get("/domains/:id", (c) => c.text("ok")); diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 720449c2..f00fba7a 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -563,6 +563,12 @@ importers: specifier: workspace:* version: link:../x devDependencies: + esbuild: + specifier: 0.27.4 + version: 0.27.4 + miniflare: + specifier: 4.20260702.0 + version: 4.20260702.0 tsup: specifier: ^8 version: 8.5.1(jiti@2.6.1)(postcss@8.5.8)(typescript@5.9.3)(yaml@2.9.0) From 16d21da03f25e380552041d4fc326b9d7581996a Mon Sep 17 00:00:00 2001 From: Rhys Sullivan <39114868+RhysSullivan@users.noreply.github.com> Date: Thu, 8 Oct 2026 02:29:15 -0700 Subject: [PATCH 5/6] Test that failures reach Analytics Engine only and observability stays off --- .../cloudflare/src/__tests__/runtime.test.ts | 66 ++++-- .../cloudflare/src/__tests__/worker.test.ts | 210 ++++++++++++------ 2 files changed, 191 insertions(+), 85 deletions(-) diff --git a/packages/@emulators/cloudflare/src/__tests__/runtime.test.ts b/packages/@emulators/cloudflare/src/__tests__/runtime.test.ts index e7a49a67..a19f697b 100644 --- a/packages/@emulators/cloudflare/src/__tests__/runtime.test.ts +++ b/packages/@emulators/cloudflare/src/__tests__/runtime.test.ts @@ -6,20 +6,25 @@ import { Log, LogLevel, Miniflare } from "miniflare"; import { afterAll, beforeAll, describe, expect, it } from "vitest"; // Runs the real Worker and Durable Object in workerd and records everything -// Cloudflare would: every tail event (the source of Workers Logs, including -// uncaught exception events) and workerd's own stdout and stderr. Faults are -// injected at the storage and stub boundaries with synthetic secrets in the -// message, the error name and the instance URL; none may reach any of it. +// Cloudflare could: every tail event (the source of Workers Logs and Issues, +// including uncaught exception events), workerd's own stdout and stderr, and +// the Analytics Engine data points the code writes. Faults are injected at the +// storage and stub boundaries with synthetic secrets in the message, the error +// name and the instance URL; none may reach any of it. const SECRET = "SYNTHETICtoken482913"; const SUFFIX = "0123456789abcdef01234567"; const INSTANCE = `synthetic-${SUFFIX}`; const SECRETS = [SECRET, SUFFIX]; +const POINT = "probe:analytics-engine"; // The probe module wraps the shipped exports. Its Durable Object hands the real // one a storage proxy that fails the way `x-probe-fault` asks, and its Worker -// can swap in a namespace whose addressing or stub fails. +// can swap in a namespace whose addressing or stub fails. Both get a FAILURES +// dataset that reports each data point as a tagged log line, so the tail sees +// exactly what Analytics Engine would store. const PROBE = ` import worker, { EmulatorDurableObject } from "./worker.ts"; +const FAILURES = { writeDataPoint: (point) => console.log("${POINT}", JSON.stringify(point)) }; const secretError = (flags) => Object.assign(new Error("uncaught ${SECRET} https://resend.${INSTANCE}.emulators.dev/emails"), { name: "${SECRET}" }, flags); export class ProbeObject extends EmulatorDurableObject { @@ -34,7 +39,7 @@ export class ProbeObject extends EmulatorDurableObject { return typeof value === "function" ? value.bind(target) : value; }, }); - super({ storage, blockConcurrencyWhile: state.blockConcurrencyWhile.bind(state) }, env); + super({ storage, blockConcurrencyWhile: state.blockConcurrencyWhile.bind(state) }, { ...env, FAILURES }); this.setFault = (next) => { fault = next; }; } async fetch(request) { @@ -44,13 +49,14 @@ export class ProbeObject extends EmulatorDurableObject { } export default { fetch(request, env) { + env = { ...env, FAILURES }; const fault = request.headers.get("x-probe-fault"); if (fault === "worker-id") env = { ...env, EMULATOR: { idFromName() { throw secretError({}); }, get: () => env.EMULATOR.get() } }; if (fault === "worker-stub") env = { ...env, EMULATOR: { idFromName: (n) => n, get: () => ({ fetch: async () => { throw secretError({ retryable: true, overloaded: true }); } }) } }; if (fault === "worker-env") - env = { get EMULATE_HOST_SUFFIX() { throw secretError({}); } }; + env = { FAILURES, get EMULATE_HOST_SUFFIX() { throw secretError({}); } }; return worker.fetch(request, env); }, }; @@ -58,7 +64,7 @@ export default { interface TailEvent { outcome: string; - event?: unknown; + event?: { request?: { url: string; headers: Record } }; entrypoint?: string; exceptions: Array<{ name: string; message: string; stack?: string }>; logs: Array<{ level: string; message: unknown[] }>; @@ -204,17 +210,37 @@ describe("emulate-hosts in workerd", () => { expect(tailed.length).toBeGreaterThanOrEqual(CASES.length); expect(tailed.map((event) => event.outcome).filter((outcome) => outcome !== "ok")).toEqual([]); expect(tailed.flatMap((event) => event.exceptions)).toEqual([]); - // Each failure is logged once, as its report. - const logged = tailed.flatMap((event) => event.logs.map((log) => log.message.map(String).join(" "))); - expect(logged).toHaveLength(CASES.length); - // Every event also carries its invocation's request (URL and headers), the - // metadata Workers Logs keeps only as invocation logs, which are off. The - // rest of each event is what the code under test can put there. - const recorded = [ - JSON.stringify(tailed.map(({ event: _event, ...rest }) => rest)), - runtimeOutput.join(""), - ...responses, - ].join("\n"); - expect(SECRETS.filter((secret) => recorded.includes(secret))).toEqual([]); + // The code writes nothing to the console: every log line is a data point + // from the probe's dataset, one per failure. + const logged = tailed.flatMap((event) => event.logs.map((log) => log.message.map(String))); + expect(logged.filter(([tag]) => tag !== POINT)).toEqual([]); + const points = logged.map(([, point]) => JSON.parse(point) as { blobs: string[] }); + expect(points).toHaveLength(CASES.length); + expect(points.map((point) => point.blobs[0]).sort()).toEqual(CASES.map((c) => c.report.error).sort()); + + // The secret is nowhere: not in the complete tail events, the runtime's + // output, the data points or the responses. + const everything = [JSON.stringify(tailed), runtimeOutput.join(""), ...responses].join("\n"); + expect(everything).not.toContain(SECRET); + // The instance name is only in the platform's own invocation metadata: the + // Worker's request URL and the object's routing headers. That metadata is + // what Workers Logs and Issues store with each record, which is why + // observability is off (see the telemetry settings test in worker.test.ts). + const carriers = new Set(); + const visit = (value: unknown, path: string) => { + if (typeof value === "string") { + if (value.includes(SUFFIX)) carriers.add(path.replace(/^\d+\./, "")); + } else if (value && typeof value === "object") { + for (const [key, child] of Object.entries(value)) visit(child, `${path}.${key}`); + } + }; + tailed.forEach((event, i) => visit(event, String(i))); + expect([...carriers].sort()).toEqual([ + "event.request.headers.x-emulator-base-url", + "event.request.headers.x-emulator-instance", + "event.request.url", + ]); + expect(SECRETS.filter((secret) => [runtimeOutput.join(""), ...responses].join("\n").includes(secret))).toEqual([]); + expect(SECRETS.filter((secret) => JSON.stringify(points).includes(secret))).toEqual([]); }, 60_000); }); diff --git a/packages/@emulators/cloudflare/src/__tests__/worker.test.ts b/packages/@emulators/cloudflare/src/__tests__/worker.test.ts index 78679ab2..30af121c 100644 --- a/packages/@emulators/cloudflare/src/__tests__/worker.test.ts +++ b/packages/@emulators/cloudflare/src/__tests__/worker.test.ts @@ -1,3 +1,4 @@ +import { readFileSync } from "node:fs"; import { describe, expect, it, vi } from "vitest"; import { EmulatorDurableObject } from "../durable-object.js"; import { instanceId } from "../diagnostics.js"; @@ -280,15 +281,58 @@ const REPORT_KEYS = [ ]; const RAY = "8f1d2c3b4a5e6f70-PHX"; +// What a failure leaves on Cloudflare: Analytics Engine data points, and +// anything written to the console, which Workers Issues keeps with the +// invocation's URL. The console must stay silent. +type DataPoint = { indexes?: string[]; blobs?: string[]; doubles?: number[] }; +const recordFailures = () => { + const points: DataPoint[] = []; + const consoles = (["log", "info", "warn", "error", "debug"] as const).map((method) => + vi.spyOn(console, method).mockImplementation(() => {}), + ); + return { + sink: { writeDataPoint: (point: DataPoint) => void points.push(point) }, + points, + stored: () => JSON.stringify(points), + consoleCalls: () => consoles.flatMap((spy) => spy.mock.calls), + restore: () => consoles.forEach((spy) => spy.mockRestore()), + }; +}; + +describe("emulate-hosts telemetry settings", () => { + // Workers Issues keeps every 5xx and error log with its invocation's URL, and + // these settings are not part of a Worker version, so the config must turn + // each one off explicitly: a deploy then also undoes a dashboard change. + const config = JSON.parse( + readFileSync(new URL("../../wrangler.jsonc", import.meta.url), "utf8").replace(/^\s*\/\/.*$/gm, ""), + ) as Record; + + it("turns every part of observability off", () => { + expect(config.observability).toEqual({ + enabled: false, + logs: { enabled: false, invocation_logs: false }, + traces: { enabled: false }, + issues: { enabled: false }, + }); + expect(config).not.toHaveProperty("tail_consumers"); + expect(config).not.toHaveProperty("logpush"); + }); + + it("records failures in the Analytics Engine dataset", () => { + expect(config.analytics_engine_datasets).toEqual([{ binding: "FAILURES", dataset: "emulate_hosts_failures" }]); + }); +}); + describe("cloudflare worker durable object failures", () => { // Cloudflare raises Durable Object stub failures as exceptions carrying // `.retryable`, `.overloaded` and `.remote` flags. const doError = (message: string, flags: { retryable?: boolean; overloaded?: boolean; remote?: boolean }) => Object.assign(new Error(message), flags); - const failingEnv = (failures: Error[]) => { + const failingEnv = (failures: Error[], sink?: Env["FAILURES"]) => { const calls: string[] = []; const env: Env = { + FAILURES: sink, EMULATOR: { idFromName: (n) => n, get: () => ({ @@ -304,18 +348,10 @@ describe("cloudflare worker durable object failures", () => { return { env, calls }; }; - const captureErrors = () => { - const spy = vi.spyOn(console, "error").mockImplementation(() => {}); - return { - text: () => spy.mock.calls.map((args) => args.map(String).join(" ")).join("\n"), - restore: () => spy.mockRestore(), - }; - }; - it("answers a retryable reset once, with a report and no replay", async () => { - const logs = captureErrors(); + const failures = recordFailures(); try { - const { env, calls } = failingEnv([doError("Network connection lost.", { retryable: true })]); + const { env, calls } = failingEnv([doError("Network connection lost.", { retryable: true })], failures.sink); const response = await worker.fetch( new Request(`https://emulators.dev/resend/${SECRET_INSTANCE}/_emulate/reset`, { method: "POST", @@ -340,9 +376,17 @@ describe("cloudflare worker durable object failures", () => { remote: false, ray: RAY, }); - expect(JSON.parse(logs.text())).toEqual(report); + // Recorded once, without the instance hash; nothing goes to the console. + expect(failures.points).toEqual([ + { + indexes: ["resend"], + blobs: ["emulator_unavailable", "resend", "POST", "/_emulate/reset", "Error", RAY], + doubles: [503, 1, 0, 0], + }, + ]); + expect(failures.consoleCalls()).toEqual([]); } finally { - logs.restore(); + failures.restore(); } }); @@ -359,14 +403,14 @@ describe("cloudflare worker durable object failures", () => { expect(await response.json()).toMatchObject({ route: "/oauth2/authorize", retryable: true }); }); - it("keeps tokens, codes, addresses and the instance out of the response and the log", async () => { - const logs = captureErrors(); + it("keeps tokens, codes, addresses and the instance out of the response and the record", async () => { + const failures = recordFailures(); try { const failure = doError( `lost while serving ${SECRET_INSTANCE}: token emu_resend_SYNTHETICtoken0001 code SYNTH-CODE-482913 for pat.synthetic@example.test`, { retryable: true }, ); - const { env } = failingEnv([failure]); + const { env } = failingEnv([failure], failures.sink); const response = await worker.fetch( new Request( `https://emulators.dev/resend/${SECRET_INSTANCE}/domains/pat.synthetic@example.test?code=SYNTH-CODE-482913`, @@ -381,17 +425,19 @@ describe("cloudflare worker durable object failures", () => { expect(Object.keys(report).sort()).toEqual(REPORT_KEYS); expect(report.route).toBe("/domains/:id"); expect(leakedSecrets(text)).toEqual([]); - expect(logs.text()).not.toBe(""); - expect(leakedSecrets(logs.text())).toEqual([]); + expect(failures.points).toHaveLength(1); + expect(leakedSecrets(failures.stored())).toEqual([]); + expect(failures.stored()).not.toContain(report.instanceId as string); + expect(failures.consoleCalls()).toEqual([]); } finally { - logs.restore(); + failures.restore(); } }); it("does not echo header or path values it cannot vouch for", async () => { - const logs = captureErrors(); + const failures = recordFailures(); try { - const { env } = failingEnv([doError("Network connection lost.", { retryable: true })]); + const { env } = failingEnv([doError("Network connection lost.", { retryable: true })], failures.sink); const response = await worker.fetch( new Request(`https://emulators.dev/pat.synthetic@example.test/${SECRET_INSTANCE}/emails`, { headers: { "cf-ray": "SYNTH-CODE-482913" }, @@ -401,9 +447,12 @@ describe("cloudflare worker durable object failures", () => { const text = await response.text(); expect(JSON.parse(text)).toMatchObject({ service: "unknown", route: "unmatched", ray: null }); expect(leakedSecrets(text)).toEqual([]); - expect(leakedSecrets(logs.text())).toEqual([]); + expect(failures.points).toHaveLength(1); + expect(leakedSecrets(failures.stored())).toEqual([]); + expect(failures.stored()).not.toContain("SYNTH-CODE-482913"); + expect(failures.consoleCalls()).toEqual([]); } finally { - logs.restore(); + failures.restore(); } }); @@ -428,10 +477,11 @@ describe("cloudflare worker durable object failures", () => { }); it("reports failures to read the body or address the object instead of throwing", async () => { - const logs = captureErrors(); + const failures = recordFailures(); try { const secretError = () => new Error(`lost ${SECRET_INSTANCE} token emu_resend_SYNTHETICtoken0001`); const addressing: Env = { + FAILURES: failures.sink, EMULATOR: { idFromName: () => { throw secretError(); @@ -446,7 +496,7 @@ describe("cloudflare worker durable object failures", () => { expect(addressed.status).toBe(500); expect(await addressed.json()).toMatchObject({ error: "emulator_unavailable", route: "/emails" }); - const { env, calls } = failingEnv([]); + const { env, calls } = failingEnv([], failures.sink); const unreadable = new Request(`https://emulators.dev/resend/${SECRET_INSTANCE}/emails`, { method: "POST", body: new ReadableStream({ @@ -460,16 +510,19 @@ describe("cloudflare worker durable object failures", () => { expect(calls).toEqual([]); expect(read.status).toBe(500); expect(await read.json()).toMatchObject({ error: "emulator_unavailable", route: "/emails" }); - expect(leakedSecrets(logs.text())).toEqual([]); + expect(failures.points).toHaveLength(2); + expect(leakedSecrets(failures.stored())).toEqual([]); + expect(failures.consoleCalls()).toEqual([]); } finally { - logs.restore(); + failures.restore(); } }); it("reports any other Worker failure instead of throwing", async () => { - const logs = captureErrors(); + const failures = recordFailures(); try { const env = { + FAILURES: failures.sink, get EMULATE_HOST_SUFFIX(): string { throw new Error(`config read failed for ${SECRET_INSTANCE}`); }, @@ -489,20 +542,31 @@ describe("cloudflare worker durable object failures", () => { remote: false, ray: null, }); - expect(leakedSecrets(text + logs.text())).toEqual([]); + expect(failures.points).toEqual([ + { + indexes: ["unknown"], + blobs: ["worker_error", "unknown", "GET", "unmatched", "Error", ""], + doubles: [500, 0, 0, 0], + }, + ]); + expect(leakedSecrets(text + failures.stored())).toEqual([]); + expect(failures.consoleCalls()).toEqual([]); } finally { - logs.restore(); + failures.restore(); } }); it("reports only allowlisted error classes and methods", async () => { - const logs = captureErrors(); + const failures = recordFailures(); try { - const { env } = failingEnv([ - doError("lost", { retryable: true }), - Object.assign(doError("lost", { retryable: true }), { name: "SYNTHETICtoken482913" }), - Object.assign(new TypeError("lost"), { retryable: true }), - ]); + const { env } = failingEnv( + [ + doError("lost", { retryable: true }), + Object.assign(doError("lost", { retryable: true }), { name: "SYNTHETICtoken482913" }), + Object.assign(new TypeError("lost"), { retryable: true }), + ], + failures.sink, + ); const send = (method: string) => worker .fetch(new Request(`https://emulators.dev/resend/${SECRET_INSTANCE}/emails`, { method }), env) @@ -510,10 +574,12 @@ describe("cloudflare worker durable object failures", () => { expect(await send("SYNTHETICTOKEN")).toMatchObject({ method: "OTHER", errorClass: "Error" }); expect(await send("GET")).toMatchObject({ method: "GET", errorClass: "other" }); expect(await send("GET")).toMatchObject({ errorClass: "TypeError" }); - expect(logs.text()).not.toContain("SYNTHETICTOKEN"); - expect(leakedSecrets(logs.text())).toEqual([]); + expect(failures.points).toHaveLength(3); + expect(failures.stored()).not.toContain("SYNTHETICTOKEN"); + expect(leakedSecrets(failures.stored())).toEqual([]); + expect(failures.consoleCalls()).toEqual([]); } finally { - logs.restore(); + failures.restore(); } }); }); @@ -620,9 +686,9 @@ describe("cloudflare durable object control plane", () => { // Builds the instance, then makes the next storage write throw `failure` once. // One-shot, so a later write (the persist after every mutating request) // cannot rethrow it and hide what the router did with the first one. - const failNextWriteAfterWarmup = async (failure: unknown) => { + const failNextWriteAfterWarmup = async (failure: unknown, sink?: Env["FAILURES"]) => { const { state } = makeState(); - const durableObject = new EmulatorDurableObject(state, {}); + const durableObject = new EmulatorDurableObject(state, { FAILURES: sink }); const warm = await durableObject.fetch( new Request("https://github.my-run.emulators.dev/_emulate/manifest", { headers: idHeaders() }), ); @@ -667,13 +733,14 @@ describe("cloudflare durable object control plane", () => { }); }); - it("keeps secrets in an emulator error out of the response and the log", async () => { - const logs = vi.spyOn(console, "error").mockImplementation(() => {}); + it("keeps secrets in an emulator error out of the response and the record", async () => { + const failures = recordFailures(); try { const durableObject = await failNextWriteAfterWarmup( new Error( `write failed for ${SECRET_INSTANCE}: token emu_resend_SYNTHETICtoken0001 code SYNTH-CODE-482913 for pat.synthetic@example.test`, ), + failures.sink, ); // The error is thrown inside the router (reset runs in the control plane), // which used to answer it with its raw message. @@ -683,11 +750,17 @@ describe("cloudflare durable object control plane", () => { expect(Object.keys(JSON.parse(text)).sort()).toEqual(REPORT_KEYS); expect(JSON.parse(text)).toMatchObject({ error: "emulator_error", route: "/_emulate/reset" }); expect(leakedSecrets(text)).toEqual([]); - const logged = logs.mock.calls.map((args) => args.map(String).join(" ")).join("\n"); - expect(logged).not.toBe(""); - expect(leakedSecrets(logged)).toEqual([]); + expect(failures.points).toEqual([ + { + indexes: ["github"], + blobs: ["emulator_error", "github", "POST", "/_emulate/reset", "Error", ""], + doubles: [500, 0, 0, 0], + }, + ]); + expect(leakedSecrets(failures.stored())).toEqual([]); + expect(failures.consoleCalls()).toEqual([]); } finally { - logs.mockRestore(); + failures.restore(); } }); @@ -733,7 +806,7 @@ describe("cloudflare durable object control plane", () => { }); it("reports flagged and plain storage failures from seed and credential requests", async () => { - const logs = vi.spyOn(console, "error").mockImplementation(() => {}); + const failures = recordFailures(); try { const credentials = () => control("/_emulate/credentials", { type: "bearer-token", login: "synthetic-user" }); const seed = () => control("/_emulate/seed", { users: [{ login: "synthetic-user" }] }); @@ -745,13 +818,15 @@ describe("cloudflare durable object control plane", () => { [credentials, "/_emulate/credentials"], [seed, "/_emulate/seed"], ] as const) { - expect(await reportOf((await failNextWriteAfterWarmup(moved())).fetch(request()))).toMatchObject({ - status: 503, - report: { error: "emulator_unavailable", route, retryable: true }, - }); + expect(await reportOf((await failNextWriteAfterWarmup(moved(), failures.sink)).fetch(request()))).toMatchObject( + { + status: 503, + report: { error: "emulator_unavailable", route, retryable: true }, + }, + ); // A host failure is not the caller's mistake: it used to come back as a // 400 quoting the raw message. - const res = await (await failNextWriteAfterWarmup(plain())).fetch(request()); + const res = await (await failNextWriteAfterWarmup(plain(), failures.sink)).fetch(request()); const text = await res.text(); expect({ status: res.status, report: JSON.parse(text) }).toMatchObject({ status: 500, @@ -759,9 +834,11 @@ describe("cloudflare durable object control plane", () => { }); expect(leakedSecrets(text)).toEqual([]); } - expect(leakedSecrets(logs.mock.calls.map((args) => args.map(String).join(" ")).join("\n"))).toEqual([]); + expect(failures.points).toHaveLength(4); + expect(leakedSecrets(failures.stored())).toEqual([]); + expect(failures.consoleCalls()).toEqual([]); } finally { - logs.mockRestore(); + failures.restore(); } }); @@ -777,25 +854,27 @@ describe("cloudflare durable object control plane", () => { }); it("reports an error whose name could carry a secret as class other", async () => { - const logs = vi.spyOn(console, "error").mockImplementation(() => {}); + const failures = recordFailures(); try { const named = Object.assign(new Error("synthetic"), { name: "SYNTHETICtoken482913" }); - const durableObject = await failNextWriteAfterWarmup(named); + const durableObject = await failNextWriteAfterWarmup(named, failures.sink); const res = await durableObject.fetch(control("/_emulate/reset", {})); const text = await res.text(); expect(JSON.parse(text)).toMatchObject({ errorClass: "other" }); expect(text).not.toContain("SYNTHETICtoken482913"); - expect(logs.mock.calls.flat().map(String).join("\n")).not.toContain("SYNTHETICtoken482913"); + expect(failures.points).toHaveLength(1); + expect(failures.stored()).not.toContain("SYNTHETICtoken482913"); + expect(failures.consoleCalls()).toEqual([]); } finally { - logs.mockRestore(); + failures.restore(); } }); it("passes the object's flagged report through the Worker as a 503", async () => { - const logs = vi.spyOn(console, "error").mockImplementation(() => {}); + const failures = recordFailures(); try { - const durableObject = await failNextWriteAfterWarmup(moved()); - const env: Env = { EMULATOR: { idFromName: (n) => n, get: () => durableObject } }; + const durableObject = await failNextWriteAfterWarmup(moved(), failures.sink); + const env: Env = { FAILURES: failures.sink, EMULATOR: { idFromName: (n) => n, get: () => durableObject } }; const response = await worker.fetch( new Request("https://emulators.dev/github/my-run/_emulate/reset", { method: "POST", @@ -811,10 +890,11 @@ describe("cloudflare durable object control plane", () => { retryable: true, ray: RAY, }); - // Reported once, by the object. - expect(logs).toHaveBeenCalledTimes(1); + // Recorded once, by the object. + expect(failures.points).toHaveLength(1); + expect(failures.consoleCalls()).toEqual([]); } finally { - logs.mockRestore(); + failures.restore(); } }); From 5347c6ba2be306245397cdcd00fad3e73eb39a65 Mon Sep 17 00:00:00 2001 From: Rhys Sullivan <39114868+RhysSullivan@users.noreply.github.com> Date: Thu, 8 Oct 2026 03:15:00 -0700 Subject: [PATCH 6/6] Wait for each failure's exact Worker and Durable Object events in the runtime probe --- .../cloudflare/src/__tests__/runtime.test.ts | 33 ++++++++++++++----- 1 file changed, 24 insertions(+), 9 deletions(-) diff --git a/packages/@emulators/cloudflare/src/__tests__/runtime.test.ts b/packages/@emulators/cloudflare/src/__tests__/runtime.test.ts index a19f697b..b34e8a89 100644 --- a/packages/@emulators/cloudflare/src/__tests__/runtime.test.ts +++ b/packages/@emulators/cloudflare/src/__tests__/runtime.test.ts @@ -136,10 +136,17 @@ afterAll(async () => { async function waitForTail(count: number): Promise { const deadline = Date.now() + 5_000; while (tailed.length < count && Date.now() < deadline) await new Promise((r) => setTimeout(r, 25)); - // Late events (a second invocation, a trailing log) would land here. - await new Promise((r) => setTimeout(r, 200)); } +// A fault injected in the object's storage reaches it, so the request is one +// Worker invocation plus one object invocation; a Worker fault stops at the +// Worker. Each event is named by its entrypoint, the object's class or none. +const invocations = (fault: string) => (fault.startsWith("do-") ? ["ProbeObject", "worker"] : ["worker"]); +const invocation = (event: TailEvent) => ({ + fault: event.event?.request?.headers["x-probe-fault"], + entrypoint: event.entrypoint ?? "worker", +}); + const CASES: Array<{ fault: string; method?: string; path: string; status: number; report: Record }> = [ { @@ -203,11 +210,18 @@ describe("emulate-hosts in workerd", () => { responses.push(text); expect({ fault: c.fault, status: response.status }).toEqual({ fault: c.fault, status: c.status }); expect(JSON.parse(text)).toMatchObject(c.report); - // The Worker's invocation, plus the object's when the request reached it. - await waitForTail(before + 1); + // Exactly this request's invocations, and nothing from an earlier one. + const expected = invocations(c.fault); + await waitForTail(before + expected.length); + expect( + tailed + .slice(before) + .map(invocation) + .sort((a, b) => a.entrypoint.localeCompare(b.entrypoint)), + ).toEqual(expected.map((entrypoint) => ({ fault: c.fault, entrypoint }))); } - expect(tailed.length).toBeGreaterThanOrEqual(CASES.length); + expect(tailed).toHaveLength(CASES.flatMap((c) => invocations(c.fault)).length); expect(tailed.map((event) => event.outcome).filter((outcome) => outcome !== "ok")).toEqual([]); expect(tailed.flatMap((event) => event.exceptions)).toEqual([]); // The code writes nothing to the console: every log line is a data point @@ -222,10 +236,11 @@ describe("emulate-hosts in workerd", () => { // output, the data points or the responses. const everything = [JSON.stringify(tailed), runtimeOutput.join(""), ...responses].join("\n"); expect(everything).not.toContain(SECRET); - // The instance name is only in the platform's own invocation metadata: the - // Worker's request URL and the object's routing headers. That metadata is - // what Workers Logs and Issues store with each record, which is why - // observability is off (see the telemetry settings test in worker.test.ts). + // The instance name is only in each invocation's request: the client's URL + // to the Worker, and the x-emulator-* headers the Worker sets on its request + // to the object. Workers Logs and Issues can store an invocation's request + // with each record, which is why observability is off (see the telemetry + // settings test in worker.test.ts). const carriers = new Set(); const visit = (value: unknown, path: string) => { if (typeof value === "string") {