diff --git a/src/js/node/timers.promises.ts b/src/js/node/timers.promises.ts index 2074fb57ba1a..6d5edcd51bce 100644 --- a/src/js/node/timers.promises.ts +++ b/src/js/node/timers.promises.ts @@ -4,26 +4,10 @@ const { validateBoolean, validateAbortSignal, validateObject, validateNumber } = require("internal/validators"); const { resistStopPropagation } = require("internal/shared"); -const symbolAsyncIterator = Symbol.asyncIterator; const setImmediateGlobal = globalThis.setImmediate; const setTimeoutGlobal = globalThis.setTimeout; const setIntervalGlobal = globalThis.setInterval; -function asyncIterator({ next: nextFunction, return: returnFunction }) { - const result = {}; - if (typeof nextFunction === "function") { - result.next = nextFunction; - } - if (typeof returnFunction === "function") { - result.return = returnFunction; - } - result[symbolAsyncIterator] = function () { - return this; - }; - - return result; -} - function setTimeout(after = 1, value, options = {}) { const arguments_ = [].concat(value ?? []); try { @@ -113,60 +97,21 @@ function setImmediate(value, options = {}) { : returnValue; } -function setInterval(after = 1, value, options = {}) { - /* eslint-disable no-undefined, no-unreachable-loop, no-loop-func */ - try { - // If after is a number, but an invalid one (too big, Infinity, NaN), we only want to emit a - // warning, not throw an error. So we can't call validateNumber as that will throw if the number - // is outside of a given range. - if (typeof after != "number") { - validateNumber(after, "delay"); - } - } catch (error) { - return asyncIterator({ - next: function () { - return Promise.$reject(error); - }, - }); - } - try { - validateObject(options, "options"); - } catch (error) { - return asyncIterator({ - next: function () { - return Promise.$reject(error); - }, - }); +async function* setInterval(after = 1, value, options = {}) { + if (typeof after !== "number") { + validateNumber(after, "delay"); } + validateObject(options, "options"); const { signal, ref: reference = true } = options; - try { - validateAbortSignal(signal, "options.signal"); - } catch (error) { - return asyncIterator({ - next: function () { - return Promise.$reject(error); - }, - }); - } - try { - validateBoolean(reference, "options.ref"); - } catch (error) { - return asyncIterator({ - next: function () { - return Promise.$reject(error); - }, - }); - } + validateAbortSignal(signal, "options.signal"); + validateBoolean(reference, "options.ref"); + if (signal?.aborted) { - return asyncIterator({ - next: function () { - return Promise.$reject($makeAbortError(undefined, { cause: signal.reason })); - }, - }); + throw $makeAbortError(undefined, { cause: signal.reason }); } - let onCancel, interval; - + let onCancel; + let interval; try { let notYielded = 0; let callback; @@ -184,50 +129,25 @@ function setInterval(after = 1, value, options = {}) { onCancel = () => { clearInterval(interval); if (callback) { - callback(); + callback(Promise.$reject($makeAbortError(undefined, { cause: signal.reason }))); callback = undefined; } }; signal.addEventListener("abort", onCancel, resistStopPropagation({ __proto__: null, once: true })); } - return asyncIterator({ - next: function () { - return new Promise((resolve, reject) => { - if (!signal?.aborted) { - if (notYielded === 0) { - callback = resolve; - } else { - resolve(); - } - } else if (notYielded === 0) { - reject($makeAbortError(undefined, { cause: signal.reason })); - } else { - resolve(); - } - }).then(() => { - if (notYielded > 0) { - notYielded = notYielded - 1; - return { done: false, value: value }; - } else if (signal?.aborted) { - throw $makeAbortError(undefined, { cause: signal.reason }); - } - return { done: true }; - }); - }, - return: function () { - clearInterval(interval); - signal?.removeEventListener("abort", onCancel); - return Promise.$resolve({}); - }, - }); - } catch { - return asyncIterator({ - next: function () { - clearInterval(interval); - signal?.removeEventListener("abort", onCancel); - }, - }); + while (!signal?.aborted) { + if (notYielded === 0) { + await new Promise(resolve => (callback = resolve)); + } + for (; notYielded > 0; notYielded--) { + yield value; + } + } + throw $makeAbortError(undefined, { cause: signal?.reason }); + } finally { + clearInterval(interval); + signal?.removeEventListener("abort", onCancel); } } diff --git a/test/js/node/timers.promises/timers.promises.test.ts b/test/js/node/timers.promises/timers.promises.test.ts index 06b210416413..62606a571c2e 100644 --- a/test/js/node/timers.promises/timers.promises.test.ts +++ b/test/js/node/timers.promises/timers.promises.test.ts @@ -1,6 +1,16 @@ import { describe, expect, it } from "bun:test"; +import { bunEnv, bunExe } from "harness"; import { setImmediate, setInterval, setTimeout } from "node:timers/promises"; +const bound = (p: Promise, ms: number) => + Promise.race([ + p.then( + v => ["settled", v] as const, + e => ["rejected", e] as const, + ), + setTimeout(ms, ["TIMEOUT"] as const), + ]); + describe("setTimeout", () => { it("abort() does not emit global error", async () => { let unhandledRejectionCaught = false; @@ -103,4 +113,101 @@ describe("setInterval", () => { await iterator.return!(); } }); + + it("does not arm the interval until the iterator is first consumed", async () => { + // An async iterator that is created but never iterated must not keep the event loop alive. + // Node arms the underlying interval lazily on the first next() (async generator semantics). + const src = ` + const { setInterval } = require("node:timers/promises"); + const it = setInterval(100000); + void it; + const t = setTimeout(() => { console.log("STILL_ALIVE"); process.exit(1); }, 2000); + t.unref(); + console.log("created"); + `; + await using proc = Bun.spawn({ + cmd: [bunExe(), "-e", src], + env: bunEnv, + stdout: "pipe", + stderr: "pipe", + }); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + expect(stderr).toBe(""); + expect(stdout).toBe("created\n"); + expect(exitCode).toBe(0); + }); + + it("settles every concurrent next() call", async () => { + const it = setInterval(10, "tick"); + try { + const first = it.next(); + const second = it.next(); + const third = it.next(); + expect(await bound(first, 500)).toEqual(["settled", { done: false, value: "tick" }]); + expect(await bound(second, 500)).toEqual(["settled", { done: false, value: "tick" }]); + expect(await bound(third, 500)).toEqual(["settled", { done: false, value: "tick" }]); + } finally { + await it.return(); + } + }); + + it("return(value) resolves { value, done: true }", async () => { + const it = setInterval(10, "x"); + await it.next(); + expect(await it.return("RV")).toEqual({ value: "RV", done: true }); + expect(await it.return("again")).toEqual({ value: "again", done: true }); + }); + + it("next() after return() resolves { done: true }", async () => { + const it = setInterval(10, "x"); + await it.next(); + await it.return(); + expect(await bound(it.next(), 500)).toEqual(["settled", { value: undefined, done: true }]); + }); + + it("next() after return() does not yield buffered ticks", async () => { + const it = setInterval(1, "buf"); + await it.next(); + await it.next(); + await setTimeout(50); + await it.return(); + expect(await bound(it.next(), 500)).toEqual(["settled", { value: undefined, done: true }]); + }); + + it("second for-await over the same iterator completes immediately", async () => { + const it = setInterval(10, "y"); + let count = 0; + for await (const _ of it) { + if (++count >= 2) break; + } + expect(count).toBe(2); + let count2 = 0; + const loop = (async () => { + for await (const _ of it) count2++; + })(); + expect(await bound(loop, 500)).toEqual(["settled", undefined]); + expect(count2).toBe(0); + }); + + it("has a throw() method that rejects and closes the iterator", async () => { + const it = setInterval(10, "x"); + expect(typeof it.throw).toBe("function"); + await it.next(); + const err = new Error("boom"); + const r = await bound(it.throw(err), 500); + expect(r[0]).toBe("rejected"); + expect(r[1]).toBe(err); + expect(await bound(it.next(), 500)).toEqual(["settled", { value: undefined, done: true }]); + }); + + it("abort rejects a pending next() and subsequent next() resolves done", async () => { + const ac = new AbortController(); + const it = setInterval(100, "tick", { signal: ac.signal }); + const p = it.next(); + ac.abort(); + const r = await bound(p, 500); + expect(r[0]).toBe("rejected"); + expect((r[1] as Error).name).toBe("AbortError"); + expect(await bound(it.next(), 500)).toEqual(["settled", { value: undefined, done: true }]); + }); });