Skip to content

stream.Transform changing order of items #46765

Description

@zuozp8

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?

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');
        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));

How often does it reproduce? Is there a required condition?

always

What is the expected behavior?

output:
 header from constructor
"firstLine"
"secondLine"

What do you see instead?

output:
 header from constructor
"secondLine"
"firstLine"

Additional information

I expect Transform to always process incoming items in order, even with reckless code like in the example.

Activity

  1. changed the title [-]Transform stream changing order of items[/-] [+]stream.Transform changing order of items[/+] on Feb 22, 2023
  2. rluvaton commented on Feb 24, 2023

    @rluvaton
    Member

    You accidentally wrote the expected in the current behavior

  3. zuozp8 commented on Feb 24, 2023

    @zuozp8
    Author

    You accidentally wrote the expected in the current behavior

    I corrected it, thank you

  4. ronag commented on Feb 24, 2023

    @ronag
    Member

    Do you think you can make an even simpler example?

  5. rluvaton commented on Feb 24, 2023

    @rluvaton
    Member

    What I saw is when you remove the construct it works well

  6. ronag commented on Feb 24, 2023

    @ronag
    Member

    That's a good hint.

  7. zuozp8 commented on Feb 24, 2023

    @zuozp8
    Author

    sometimes it starts/stops occuring when you change some process.nextTick(callback) into callack() or the other way around, so it's some kind of race-condition

    I wasn't able to replicate the problem without nested Transforms

  8. rluvaton commented on Feb 24, 2023

    @rluvaton
    Member

    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

  9. rluvaton commented on Feb 24, 2023

    @rluvaton
    Member

    Found another solution and updated my PR

  10. mcollina commented on Feb 24, 2023

    @mcollina
    SponsorMember

    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 nextTick in the constructor logic.

  11. mcollina commented on Feb 24, 2023

    @mcollina
    SponsorMember
  12. added
    confirmed-bugIssues and PRs for confirmed bugs.
    good first issueIssues that are suitable for first-time contributors.
    on Feb 24, 2023
  13. rluvaton commented on Feb 24, 2023

    @rluvaton
    Member

    I think 1f75a95/lib/internal/streams/readable.js#LL222C7-L222C7 should be delayed with a nextTick.

    delaying it with nextTick did not worked...

  14. mcollina commented on Feb 24, 2023

    @mcollina
    SponsorMember

    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.

  15. rluvaton commented on Feb 24, 2023

    @rluvaton
    Member

    I tested with your simplified version

  16. mcollina commented on Feb 24, 2023

    @mcollina
    SponsorMember

    Here is the fix: #46818

  17. removed
    good first issueIssues that are suitable for first-time contributors.
    on Feb 24, 2023
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    confirmed-bugIssues and PRs for confirmed bugs.streamIssues and PRs related to Node.js streams.

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions