Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/eff-549-worker-send-error.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
"effect": patch
---

Report buffered worker send failures as `WorkerError` values.
27 changes: 13 additions & 14 deletions packages/effect/src/unstable/workers/Worker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -165,6 +165,17 @@ export const makePlatform = <W>() =>
const spawn = (yield* Spawner) as SpawnerFn<W>
let currentPort: P | undefined
const buffer: Array<[unknown, ReadonlyArray<unknown> | undefined]> = []
const sendToPort = (port: P, message: unknown, transfers?: ReadonlyArray<unknown>) =>
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 = <A, E, R>(
handler: (_: O) => Effect.Effect<A, E, R>,
Expand Down Expand Up @@ -209,7 +220,7 @@ export const makePlatform = <W>() =>
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
}
Expand All @@ -224,19 +235,7 @@ export const makePlatform = <W>() =>
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 }
Expand Down
42 changes: 41 additions & 1 deletion packages/effect/test/unstable/workers/WorkerError.test.ts
Original file line number Diff line number Diff line change
@@ -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,
Expand Down Expand Up @@ -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<void>()
let emit!: (message: Worker.PlatformMessage) => void
const platform = Worker.makePlatform<object>()({
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")
}
}
}))
})
})
Loading