Repository navigation
stream.Transform changing order of items #46765
Description
Activity
- changed the title
[-]Transform stream changing order of items[/-][+]stream.Transform changing order of items[/+]on Feb 22, 2023 - addedstreamIssues and PRs related to Node.js streams.Issues and PRs related to Node.js streams.
on Feb 23, 2023 You accidentally wrote the expected in the current behavior
You accidentally wrote the expected in the current behavior
I corrected it, thank you
Reacted by Raz LuvatonDo you think you can make an even simpler example?
What I saw is when you remove the
constructit works wellThat's a good hint.
sometimes it starts/stops occuring when you change some
process.nextTick(callback)intocallack()or the other way around, so it's some kind of race-conditionI wasn't able to replicate the problem without nested Transforms
When you wrap the:
else if (innerTranfrorm.write(row)) { process.nextTick(callback); } else
like this:
} else if (innerTransform.write('outer | ' + row)) { process.nextTick(() => { process.nextTick(callback); }); } else {
it solves the problem
Found another solution and updated my PR
This is indeed a bug. Here is a simplified version:
const stream = require("node:stream"); const s = new stream.Transform({ objectMode: true, construct(callback) { this.push('header from constructor\n'); callback(); }, transform: (row, encoding, callback) => { callback(null, JSON.stringify(row) + '\n'); }, }); s.pipe(process.stdout); s.write('firstLine'); process.nextTick(() => s.write('secondLine'));
This fixes the example:
const stream = require("node:stream"); const s = new stream.Transform({ objectMode: true, construct(callback) { this.push('header from constructor\n'); process.nextTick(callback); }, transform: (row, encoding, callback) => { callback(null, JSON.stringify(row) + '\n'); }, }); s.pipe(process.stdout); s.write('firstLine'); process.nextTick(() => s.write('secondLine'));
So I think we are missing a
nextTickin the constructor logic.I think https://github.com/nodejs/node/blob/1f75a9513fc829fbd5326a5c9bfa7678b7431eb8/lib/internal/streams/readable.js#LL222C7-L222C7 should be delayed with a
nextTick.- addedconfirmed-bugIssues and PRs for confirmed bugs.Issues and PRs for confirmed bugs.good first issueIssues that are suitable for first-time contributors.Issues that are suitable for first-time contributors.
on Feb 24, 2023 I think
1f75a95/lib/internal/streams/readable.js#LL222C7-L222C7 should be delayed with anextTick.delaying it with
nextTickdid not worked...Given the following script:
const stream = require("node:stream"); const consumers = require("node:stream/consumers"); const createInnerTransfrom = () => new stream.Transform({ objectMode: true, construct(callback) { this.push('header from constructor\n'); process.nextTick(callback); }, transform: (row, encoding, callback) => { callback(null, JSON.stringify(row) + '\n'); }, }); const createOuterTransfrom = () => { let innerTranfrorm; return new stream.Transform({ objectMode: true, transform(row, encoding, callback) { if (!innerTranfrorm) { innerTranfrorm = createInnerTransfrom(); innerTranfrorm.on('data', (data) => this.push(data)); callback(); } else if (innerTranfrorm.write(row)) { process.nextTick(callback); } else { innerTranfrorm.once('drain', callback); } }, }); }; consumers.text(stream.Readable.from([ 'create InnerTransform', 'firstLine', 'secondLine', ]).pipe(createOuterTransfrom())).then((text) => console.log('output:\n', text));It will error:
node:internal/process/promises:288 triggerUncaughtException(err, true /* fromPromise */); ^ Error [ERR_STREAM_PUSH_AFTER_EOF]: stream.push() after EOF at new NodeError (node:internal/errors:399:5) at readableAddChunk (node:internal/streams/readable:285:30) at Readable.push (node:internal/streams/readable:234:10) at Transform.<anonymous> (/Users/matteo/tmp/aaa.js:20:58) at Transform.emit (node:events:513:28) at addChunk (node:internal/streams/readable:324:12) at readableAddChunk (node:internal/streams/readable:297:9) at Readable.push (node:internal/streams/readable:234:10) at node:internal/streams/transform:182:12 at Transform.transform [as _transform] (/Users/matteo/tmp/aaa.js:10:9) { code: 'ERR_STREAM_PUSH_AFTER_EOF' }This is correct, as you are pushing things after the stream ended.
I tested with your simplified version
Here is the fix: #46818
- added a commit that references this issue
on Feb 24, 2023 - removedgood first issueIssues that are suitable for first-time contributors.Issues that are suitable for first-time contributors.
on Feb 24, 2023 - added a commit that references this issue
on Feb 26, 2023 - added 2 commits that reference this issue
on Mar 13, 2023 - added a commit that references this issue
on Apr 11, 2023
Version
v18.14.1
Platform
Linux 703748a51615 5.18.8-051808-generic #202206290850 SMP PREEMPT_DYNAMIC Wed Jun 29 08:59:08 UTC 2022 x86_64 Linux
Subsystem
stream
What steps will reproduce the bug?
How often does it reproduce? Is there a required condition?
always
What is the expected behavior?
What do you see instead?
Additional information
I expect
Transformto always process incoming items in order, even with reckless code like in the example.