diff --git a/src/worker-entry.ts b/src/worker-entry.ts index 22594c4..95ace5c 100644 --- a/src/worker-entry.ts +++ b/src/worker-entry.ts @@ -14,6 +14,10 @@ function isTaskArg(arg: unknown): arg is { __task__: number; id: string; args: u return typeof arg === 'object' && arg !== null && '__task__' in arg; } +function needsAsyncResolve(arg: unknown): boolean { + return arg instanceof MessagePort || isTaskArg(arg); +} + // Returns the function synchronously when cached — callers must not await // unconditionally, or they pay a microtask hop for every dispatch. function resolveFn(id: string): Fn | Promise { @@ -101,29 +105,6 @@ function portToAsyncIterable(port: MessagePort): AsyncIterable { }; } -// Builds the resolved args array imperatively. Returns a plain array when -// every arg can be resolved synchronously; only switches to a Promise at -// the first arg that needs async resolution. This matters because awaiting -// a plain array still costs a microtask hop. -function resolveArgs(args: unknown[]): unknown[] | Promise { - const result: unknown[] = new Array(args.length); - for (let i = 0; i < args.length; i++) { - const arg = args[i]; - if (arg instanceof MessagePort || isTaskArg(arg)) { - return resolveArgsAsync(args, result, i); - } - result[i] = deserializeArg(arg); - } - return result; -} - -async function resolveArgsAsync(args: unknown[], partial: unknown[], from: number): Promise { - for (let i = from; i < args.length; i++) { - partial[i] = await resolveArg(args[i]); - } - return partial; -} - async function resolveArg(arg: unknown): Promise { if (arg instanceof MessagePort) { return portToAsyncIterable(arg); @@ -132,8 +113,12 @@ async function resolveArg(arg: unknown): Promise { if (taskCache.has(arg.__task__)) { return taskCache.get(arg.__task__); } - const argsM = resolveArgs(arg.args); - const resolvedArgs = argsM instanceof Promise ? await argsM : argsM; + const inputArgs = arg.args; + const resolvedArgs = new Array(inputArgs.length); + for (let i = 0; i < inputArgs.length; i++) { + const a = inputArgs[i]; + resolvedArgs[i] = needsAsyncResolve(a) ? await resolveArg(a) : deserializeArg(a); + } const fnM = resolveFn(arg.id); const fn = fnM instanceof Promise ? await fnM : fnM; const raw = fn(...resolvedArgs); @@ -181,15 +166,27 @@ function handleValueTask(msg: TaskMsg): void { function invokeWithArgs(callId: number, fn: Fn, args: unknown[]): void { try { - const argsM = resolveArgs(args); - if (argsM instanceof Promise) { - argsM.then( - (resolved) => invokeAndRespond(callId, fn, resolved), - (err) => postError(callId, err), - ); + // Sync prologue: walk args in-place until we hit one that needs async + // resolution. If we finish the loop synchronously, go straight to invoke. + const resolved = new Array(args.length); + let i = 0; + for (; i < args.length; i++) { + const arg = args[i]; + if (needsAsyncResolve(arg)) break; + resolved[i] = deserializeArg(arg); + } + if (i === args.length) { + invokeAndRespond(callId, fn, resolved); return; } - invokeAndRespond(callId, fn, argsM); + // Async tail: same loop body, awaiting where needed. + (async () => { + for (; i < args.length; i++) { + const arg = args[i]; + resolved[i] = needsAsyncResolve(arg) ? await resolveArg(arg) : deserializeArg(arg); + } + invokeAndRespond(callId, fn, resolved); + })().catch((err) => postError(callId, err)); } catch (err) { postError(callId, err); } @@ -217,8 +214,11 @@ async function handleStreamTask(msg: TaskMsg): Promise { try { const fnM = resolveFn(id); const fn = fnM instanceof Promise ? await fnM : fnM; - const argsM = resolveArgs(args); - const resolvedArgs = argsM instanceof Promise ? await argsM : argsM; + const resolvedArgs = new Array(args.length); + for (let i = 0; i < args.length; i++) { + const a = args[i]; + resolvedArgs[i] = needsAsyncResolve(a) ? await resolveArg(a) : deserializeArg(a); + } let paused = false; let resumed: (() => void) | null = null; @@ -244,6 +244,20 @@ async function handleStreamTask(msg: TaskMsg): Promise { } }); + // Adaptive yield: arm a setImmediate whose callback flips a flag. + // If the generator's natural awaits tick the event loop between emits, + // the flag flips on its own and we never force a yield — parent-side + // pause/close messages arrive via the loop's normal message delivery. + // If the generator is purely synchronous (no natural ticks), the flag + // stays unflipped and once `streak` hits YIELD_EVERY we force a yield + // so pause/close can actually reach us. + const YIELD_EVERY = 16; + let ticked = false; + setImmediate(() => { + ticked = true; + }); + let streak = 0; + try { const gen = fn(...resolvedArgs) as AsyncGenerator; for await (const value of gen) { @@ -257,7 +271,21 @@ async function handleStreamTask(msg: TaskMsg): Promise { const transferList: Transferable[] = []; collectTransferables(value, transferList); port!.postMessage({ value: serialized, done: false }, transferList); - await new Promise((r) => setImmediate(r)); + streak++; + if (ticked) { + ticked = false; + setImmediate(() => { + ticked = true; + }); + streak = 0; + } else if (streak >= YIELD_EVERY) { + await new Promise((r) => setImmediate(r)); + ticked = false; + setImmediate(() => { + ticked = true; + }); + streak = 0; + } } if (!cancelled) port!.postMessage({ done: true }); } catch (err) {