Skip to content
Open
Show file tree
Hide file tree
Changes from 1 commit
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
128 changes: 25 additions & 103 deletions src/js/node/timers.promises.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -113,60 +97,23 @@ 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 after is a number, but an invalid one (too big, Infinity, NaN), we only want to emit a
// warning, not throw an error, so only validate non-number inputs here.
Comment thread
robobun marked this conversation as resolved.
Outdated
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;
Expand All @@ -184,50 +131,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);
}
}

Expand Down
107 changes: 107 additions & 0 deletions test/js/node/timers.promises/timers.promises.test.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,15 @@
import { describe, expect, it } from "bun:test";
import { setImmediate, setInterval, setTimeout } from "node:timers/promises";
import { bunEnv, bunExe } from "harness";

const bound = <T>(p: Promise<T>, 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 () => {
Expand Down Expand Up @@ -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();
}
});
Comment thread
robobun marked this conversation as resolved.

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 }]);
});
Comment thread
robobun marked this conversation as resolved.

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 }]);
});
});
Loading