diff --git a/.changeset/eff-549-worker-send-error.md b/.changeset/eff-549-worker-send-error.md new file mode 100644 index 00000000000..5bab011f5f4 --- /dev/null +++ b/.changeset/eff-549-worker-send-error.md @@ -0,0 +1,5 @@ +--- +"effect": patch +--- + +Report buffered worker send failures as `WorkerError` values. diff --git a/packages/effect/src/unstable/workers/Worker.ts b/packages/effect/src/unstable/workers/Worker.ts index c2c65acf503..c96b6586310 100644 --- a/packages/effect/src/unstable/workers/Worker.ts +++ b/packages/effect/src/unstable/workers/Worker.ts @@ -165,6 +165,17 @@ export const makePlatform = () => const spawn = (yield* Spawner) as SpawnerFn let currentPort: P | undefined const buffer: Array<[unknown, ReadonlyArray | undefined]> = [] + const sendToPort = (port: P, message: unknown, transfers?: ReadonlyArray) => + Effect.try({ + try: () => port.postMessage([0, message], transfers as any), + catch: (cause) => + new WorkerError({ + reason: new WorkerSendError({ + message: "Failed to send message to worker", + cause + }) + }) + }) const run = ( handler: (_: O) => Effect.Effect, @@ -209,7 +220,7 @@ export const makePlatform = () => currentPort = port if (buffer.length > 0) { for (const [message, transfers] of buffer) { - port.postMessage([0, message], transfers as any) + yield* sendToPort(port, message, transfers) } buffer.length = 0 } @@ -224,19 +235,7 @@ export const makePlatform = () => buffer.push([message, transfers]) return Effect.void } - try { - currentPort.postMessage([0, message], transfers as any) - return Effect.void - } catch (cause) { - return Effect.fail( - new WorkerError({ - reason: new WorkerSendError({ - message: "Failed to send message to worker", - cause - }) - }) - ) - } + return sendToPort(currentPort, message, transfers) }) return { run, send } diff --git a/packages/effect/test/unstable/workers/WorkerError.test.ts b/packages/effect/test/unstable/workers/WorkerError.test.ts index 7088cf930ec..9ef0b368243 100644 --- a/packages/effect/test/unstable/workers/WorkerError.test.ts +++ b/packages/effect/test/unstable/workers/WorkerError.test.ts @@ -1,5 +1,6 @@ import { assert, describe, it } from "@effect/vitest" -import { Effect, Schema } from "effect" +import { Cause, Deferred, Effect, Exit, Fiber, Schema } from "effect" +import * as Worker from "effect/unstable/workers/Worker" import { isWorkerError, WorkerError, @@ -87,4 +88,43 @@ describe("WorkerError", () => { assert.strictEqual(decoded.message, "Failed to send message") })) }) + + describe("buffered sends", () => { + it.effect("reports postMessage failures as WorkerSendError", () => + Effect.gen(function*() { + const listening = yield* Deferred.make() + let emit!: (message: Worker.PlatformMessage) => void + const platform = Worker.makePlatform()({ + setup: () => + Effect.succeed({ + postMessage() { + throw new Error("post failed") + } + }), + listen: (options) => + Effect.sync(() => { + emit = options.emit + }).pipe(Effect.andThen(Deferred.succeed(listening, void 0))) + }) + const worker = yield* platform.spawn(0).pipe( + Effect.provideService(Worker.Spawner, () => ({})) + ) + + yield* worker.send("buffered") + const run = yield* Effect.forkChild(worker.run(() => Effect.void)) + yield* Deferred.await(listening) + yield* Effect.sync(() => emit([0])) + const exit = yield* Fiber.await(run) + + assert(Exit.isFailure(exit)) + if (Exit.isFailure(exit)) { + assert.isFalse(Cause.hasDies(exit.cause)) + const error = Cause.squash(exit.cause) + assert.isTrue(isWorkerError(error)) + if (isWorkerError(error)) { + assert.strictEqual(error.reason._tag, "WorkerSendError") + } + } + })) + }) })