From 51372010abe1660c74b15fe12a76e8b38236c8e3 Mon Sep 17 00:00:00 2001 From: Devin Ivy Date: Wed, 15 Apr 2026 10:21:50 -0400 Subject: [PATCH] feat: preserve error details across worker boundary Errors thrown in moroutines now transfer via structured clone instead of extracting message strings. This preserves stack traces pointing to the actual source, built-in subclass identity (TypeError, RangeError, etc.), and cause chains. Applies to regular dispatch, pool dispatch, and streaming tasks. Co-Authored-By: Claude Opus 4.6 (1M context) --- .changeset/error-transfer.md | 7 +++++ src/execute.ts | 11 +++---- src/worker-entry.ts | 10 +++--- test/error.test.ts | 61 +++++++++++++++++++++++++++++++++++- test/fixtures/math.ts | 8 +++++ test/fixtures/stream-gen.ts | 7 +++++ 6 files changed, 91 insertions(+), 13 deletions(-) create mode 100644 .changeset/error-transfer.md diff --git a/.changeset/error-transfer.md b/.changeset/error-transfer.md new file mode 100644 index 0000000..9f8b194 --- /dev/null +++ b/.changeset/error-transfer.md @@ -0,0 +1,7 @@ +--- +'moroutine': minor +--- + +Preserve error details across worker boundary + +Errors thrown in moroutines now transfer with `message`, `stack`, `cause`, and built-in subclass identity (`TypeError`, `RangeError`, etc.) preserved via structured clone. Previously only the message string was kept. diff --git a/src/execute.ts b/src/execute.ts index b1ca026..6cf1d32 100644 --- a/src/execute.ts +++ b/src/execute.ts @@ -16,12 +16,12 @@ const pending = new Map void; reject: (reason const streamPortStack: MessagePort[][] = []; export function setupWorker(worker: Worker): void { - worker.on('message', (msg: { callId: number; value?: unknown; error?: string }) => { + worker.on('message', (msg: { callId: number; value?: unknown; error?: Error }) => { const call = pending.get(msg.callId); if (!call) return; pending.delete(msg.callId); if (msg.error !== undefined) { - call.reject(new Error(msg.error)); + call.reject(msg.error); } else { call.resolve(deserializeArg(msg.value)); } @@ -74,8 +74,7 @@ function pipeToPort(iterable: AsyncIterable, port: MessagePort, highWat if (!cancelled) port.postMessage({ done: true }); } catch (err) { if (!cancelled) { - const message = err instanceof Error ? err.message : String(err); - port.postMessage({ done: true, error: message }); + port.postMessage({ done: true, error: err instanceof Error ? err : new Error(String(err)) }); } } try { @@ -174,9 +173,9 @@ export function dispatchStream( let paused = false; let waiting: (() => void) | null = null; - port1.on('message', (msg: { value?: unknown; done?: boolean; error?: string }) => { + port1.on('message', (msg: { value?: unknown; done?: boolean; error?: Error }) => { if (msg.error) { - error = new Error(msg.error); + error = msg.error; done = true; port1.close(); resolveDone!(); diff --git a/src/worker-entry.ts b/src/worker-entry.ts index 5f0b5b9..4e5b6b1 100644 --- a/src/worker-entry.ts +++ b/src/worker-entry.ts @@ -20,9 +20,9 @@ function portToAsyncIterable(port: MessagePort): AsyncIterable { const HIGH_WATER = 16; - port.on('message', (msg: { value?: unknown; done?: boolean; error?: string }) => { + port.on('message', (msg: { value?: unknown; done?: boolean; error?: Error }) => { if (msg.error) { - error = new Error(msg.error); + error = msg.error; done = true; if (waiting) { waiting(); @@ -160,8 +160,7 @@ parentPort!.on('message', async (msg: { callId?: number; id: string; args: unkno if (!cancelled) port.postMessage({ done: true }); } catch (err) { if (!cancelled) { - const message = err instanceof Error ? err.message : String(err); - port.postMessage({ done: true, error: message }); + port.postMessage({ done: true, error: err instanceof Error ? err : new Error(String(err)) }); } } try { @@ -177,7 +176,6 @@ parentPort!.on('message', async (msg: { callId?: number; id: string; args: unkno collectTransferables(value, transferList); parentPort!.postMessage({ callId, value: returnValue }, transferList); } catch (err) { - const message = err instanceof Error ? err.message : String(err); - parentPort!.postMessage({ callId: msg.callId!, error: message }); + parentPort!.postMessage({ callId: msg.callId!, error: err instanceof Error ? err : new Error(String(err)) }); } }); diff --git a/test/error.test.ts b/test/error.test.ts index ef769df..f83203c 100644 --- a/test/error.test.ts +++ b/test/error.test.ts @@ -1,7 +1,8 @@ import { describe, it } from 'node:test'; import assert from 'node:assert/strict'; import { workers } from 'moroutine'; -import { fail } from './fixtures/math.ts'; +import { fail, failType, failCause } from './fixtures/math.ts'; +import { failAfterType } from './fixtures/stream-gen.ts'; describe('error handling', () => { it('rejects with error from dedicated worker', async () => { @@ -16,4 +17,62 @@ describe('error handling', () => { message: 'pool boom', }); }); + + it('preserves stack trace pointing to worker source', async () => { + await using run = workers(1); + try { + await run(fail('stack check')); + assert.fail('should have thrown'); + } catch (err) { + assert.ok(err instanceof Error); + assert.match(err.stack!, /fixtures\/math\.ts/); + } + }); + + it('preserves built-in error subclass identity', async () => { + await using run = workers(1); + try { + await run(failType('type check')); + assert.fail('should have thrown'); + } catch (err) { + assert.ok(err instanceof TypeError); + assert.equal(err.message, 'type check'); + } + }); + + it('preserves error cause', async () => { + await using run = workers(1); + try { + await run(failCause('with cause')); + assert.fail('should have thrown'); + } catch (err) { + assert.ok(err instanceof Error); + assert.ok(err.cause instanceof RangeError); + assert.equal((err.cause as Error).message, 'root cause'); + } + }); + + it('preserves error details on dedicated worker', async () => { + try { + await failType('dedicated type'); + assert.fail('should have thrown'); + } catch (err) { + assert.ok(err instanceof TypeError); + assert.match((err as Error).stack!, /fixtures\/math\.ts/); + } + }); + + it('preserves error details on streaming task', async () => { + await using run = workers(1); + try { + for await (const _ of run(failAfterType(2))) { + // consume yields until error + } + assert.fail('should have thrown'); + } catch (err) { + assert.ok(err instanceof TypeError); + assert.equal((err as Error).message, 'stream type error'); + assert.match((err as Error).stack!, /fixtures\/stream-gen\.ts/); + } + }); }); diff --git a/test/fixtures/math.ts b/test/fixtures/math.ts index 2ff285d..c0ff3b0 100644 --- a/test/fixtures/math.ts +++ b/test/fixtures/math.ts @@ -5,3 +5,11 @@ export const add = mo(import.meta, (a: number, b: number) => a + b); export const fail = mo(import.meta, (msg: string) => { throw new Error(msg); }); + +export const failType = mo(import.meta, (msg: string) => { + throw new TypeError(msg); +}); + +export const failCause = mo(import.meta, (msg: string) => { + throw new Error(msg, { cause: new RangeError('root cause') }); +}); diff --git a/test/fixtures/stream-gen.ts b/test/fixtures/stream-gen.ts index 7160879..ee5ba7c 100644 --- a/test/fixtures/stream-gen.ts +++ b/test/fixtures/stream-gen.ts @@ -12,3 +12,10 @@ export const failAfter = mo(import.meta, async function* (n: number) { } throw new Error('intentional error'); }); + +export const failAfterType = mo(import.meta, async function* (n: number) { + for (let i = 0; i < n; i++) { + yield i; + } + throw new TypeError('stream type error'); +}); -- 2.51.2