Skip to content
Open
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
42 changes: 0 additions & 42 deletions src/js/internal/stream.ts
Original file line number Diff line number Diff line change
@@ -1,10 +1,8 @@
"use strict";

const ObjectKeys = Object.keys;
const ObjectDefineProperty = Object.defineProperty;

const customPromisify = Symbol.for("nodejs.util.promisify.custom");
const { streamReturningOperators, promiseReturningOperators } = require("internal/streams/operators");
const compose = require("internal/streams/compose");
const { setDefaultHighWaterMark, getDefaultHighWaterMark } = require("internal/streams/state");
const { pipeline } = require("internal/streams/pipeline");
Expand All @@ -22,46 +20,6 @@ Stream.isReadable = utils.isReadable;
Stream.isWritable = utils.isWritable;

Stream.Readable = require("internal/streams/readable");
const streamKeys = ObjectKeys(streamReturningOperators);
for (let i = 0; i < streamKeys.length; i++) {
const key = streamKeys[i];
const op = streamReturningOperators[key];
function fn(...args) {
if (new.target) {
throw $ERR_ILLEGAL_CONSTRUCTOR();
}
return Stream.Readable.from(op.$apply(this, args));
}
ObjectDefineProperty(fn, "name", { __proto__: null, value: op.name });
ObjectDefineProperty(fn, "length", { __proto__: null, value: op.length });
ObjectDefineProperty(Stream.Readable.prototype, key, {
__proto__: null,
value: fn,
enumerable: false,
configurable: true,
writable: true,
});
}
const promiseKeys = ObjectKeys(promiseReturningOperators);
for (let i = 0; i < promiseKeys.length; i++) {
const key = promiseKeys[i];
const op = promiseReturningOperators[key];
function fn(...args) {
if (new.target) {
throw $ERR_ILLEGAL_CONSTRUCTOR();
}
return Promise.$resolve().then(() => op.$apply(this, args));
}
ObjectDefineProperty(fn, "name", { __proto__: null, value: op.name });
ObjectDefineProperty(fn, "length", { __proto__: null, value: op.length });
ObjectDefineProperty(Stream.Readable.prototype, key, {
__proto__: null,
value: fn,
enumerable: false,
configurable: true,
writable: true,
});
}
Stream.Writable = require("internal/streams/writable");
Stream.Duplex = require("internal/streams/duplex");
Stream.Transform = require("internal/streams/transform");
Expand Down
44 changes: 44 additions & 0 deletions src/js/internal/streams/readable.ts
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@ const { SafeSet } = require("internal/primordials");
const { kAutoDestroyed } = require("internal/shared");

const ObjectDefineProperties = Object.defineProperties;
const ObjectDefineProperty = Object.defineProperty;
const SymbolAsyncDispose = Symbol.asyncDispose;
const NumberIsNaN = Number.isNaN;
const NumberIsInteger = Number.isInteger;
Expand Down Expand Up @@ -1724,4 +1725,47 @@ Readable.wrap = function (src, options) {
}).wrap(src);
};

// Node does this in lib/stream.js; here not every user of Readable loads node:stream.
const { streamReturningOperators, promiseReturningOperators } = require("internal/streams/operators");
const opStreamKeys = ObjectKeys(streamReturningOperators);
for (let i = 0; i < opStreamKeys.length; i++) {
const key = opStreamKeys[i];
const op = streamReturningOperators[key];
function fn(...args) {
if (new.target) {
throw $ERR_ILLEGAL_CONSTRUCTOR();
}
return Readable.from(op.$apply(this, args));
}
ObjectDefineProperty(fn, "name", { __proto__: null, value: op.name });
ObjectDefineProperty(fn, "length", { __proto__: null, value: op.length });
ObjectDefineProperty(Readable.prototype, key, {
__proto__: null,
value: fn,
enumerable: false,
configurable: true,
writable: true,
});
}
const opPromiseKeys = ObjectKeys(promiseReturningOperators);
for (let i = 0; i < opPromiseKeys.length; i++) {
const key = opPromiseKeys[i];
const op = promiseReturningOperators[key];
function fn(...args) {
if (new.target) {
throw $ERR_ILLEGAL_CONSTRUCTOR();
}
return Promise.$resolve().then(() => op.$apply(this, args));
}
ObjectDefineProperty(fn, "name", { __proto__: null, value: op.name });
ObjectDefineProperty(fn, "length", { __proto__: null, value: op.length });
ObjectDefineProperty(Readable.prototype, key, {
__proto__: null,
value: fn,
enumerable: false,
configurable: true,
writable: true,
});
}

export default Readable as unknown as typeof import("node:stream").Readable;
82 changes: 82 additions & 0 deletions test/js/node/stream/node-stream.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -1745,6 +1745,88 @@ describe("stream operators argument validation (nodejs/node#59529)", () => {
});
});

// These modules reach the Readable class without loading node:stream, which is where the
// operators used to be installed. Every script below runs in a fresh process for that reason.
describe("stream operators exist on Readables created without loading node:stream", () => {
const operators = [
"drop",
"filter",
"flatMap",
"map",
"take",
"every",
"forEach",
"reduce",
"toArray",
"some",
"find",
];

async function run(script) {
// `bun -e` exposes the builtin modules as globals, and a top-level `const stream = ...` in the
// script would be enough to load node:stream. Block scoping the script keeps it away from them.
await using proc = Bun.spawn({
cmd: [bunExe(), "-e", `{${script}}`],
env: bunEnv,
stdout: "pipe",
stderr: "pipe",
});
const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]);
return { stdout, stderr, exitCode };
}

it.concurrent.each([
["net.Socket", `new (require("node:net").Socket)()`],
["crypto.createHash()", `require("node:crypto").createHash("sha256")`],
[
"child_process stdout",
`require("node:child_process").spawn(process.execPath, ["--version"], { stdio: ["ignore", "pipe", "ignore"] }).stdout`,
],
["_stream_readable", `require("node:_stream_readable").prototype`],
["_stream_duplex", `require("node:_stream_duplex").prototype`],
["_stream_transform", `require("node:_stream_transform").prototype`],
["_stream_passthrough", `require("node:_stream_passthrough").prototype`],
])("%s", async (_, expression) => {
const result = await run(`
const target = ${expression};
const found = {};
for (const name of ${JSON.stringify(operators)}) found[name] = typeof target[name];
console.log(JSON.stringify(found));
`);
expect(result).toEqual({
stdout: JSON.stringify(Object.fromEntries(operators.map(name => [name, "function"]))) + "\n",
stderr: "",
exitCode: 0,
});
});

it.concurrent("map() and toArray() work on a net.Socket", async () => {
const result = await run(`
const net = require("node:net");
const server = net.createServer(socket => socket.end("hello"));
server.listen(0, "127.0.0.1", async () => {
const socket = net.connect(server.address().port, "127.0.0.1");
const chunks = await socket.map(chunk => chunk.toString().toUpperCase()).toArray();
server.close();
console.log(chunks.join(""));
});
`);
expect(result).toEqual({ stdout: "HELLO\n", stderr: "", exitCode: 0 });
});

it.concurrent("filter(), map() and toArray() work on a _stream_readable Readable", async () => {
const result = await run(`
const Readable = require("node:_stream_readable");
Readable.from([1, 2, 3, 4, 5])
.filter(n => n % 2 === 1)
.map(n => n * 10)
.toArray()
.then(values => console.log(JSON.stringify(values)));
`);
expect(result).toEqual({ stdout: "[10,30,50]\n", stderr: "", exitCode: 0 });
});
});

describe("duplexPair teardown (test-duplex-error.js)", () => {
const once = (emitter, event) => new Promise(resolve => emitter.once(event, resolve));

Expand Down
Loading