Skip to content

streams: Add support for non static compose #43664

Description

@rluvaton

What is the problem this feature will solve?

It will make the code much easier to read/write when the current readable operators are not enough

What is the feature you are proposing to solve the problem?

Adding a non-static compose method so we can chain multiple operators together

it would like like this for example:

fs.createReadStream("./file.pcapng")
  .compose(async function* (source) {
    let block = Buffer.alloc(0);

    for await (const chunk of source) {
      block = Buffer.concat(block, chunk);

      if (block.length < REQUIRED_BLOCK_SIZE_TO_GET_TO_TOTAL_LENGTH) {
        continue;
      }

      const blockExpectedTotalSize = getTotalSizeFromBlock(block);

      if (block.length < blockExpectedTotalSize) {
        continue;
      }

      if (mergedBlock.length >= blockExpectedTotalSize) {
        yield block.slice(0, blockExpectedTotalSize);
        block = block.slice(blockExpectedTotalSize);
      }
    }
  })
  .map((blockBuffer) => parseBufferBasedOnBlockType(blockBuffer))
  .filter((parsedBlock) => parsedblock.dest === "192.168.1.1")
  .toArray();

What alternatives have you considered?

Using static compose

Activity

  1. rluvaton commented on Jul 3, 2022

    @rluvaton
    MemberAuthor
  2. benjamingr commented on Jul 3, 2022

    @benjamingr
    Member

    @nodejs/streams

  3. mcollina commented on Jul 3, 2022

    @mcollina
    SponsorMember

    What use cases would this solve vs the static compose?

  4. rluvaton commented on Jul 3, 2022

    @rluvaton
    MemberAuthor

    no real problem but make the code much more readable and consistent with the chaining API

    Let's say you have a transformer at the start and at the middle, it will be really hard to read:

    Example using static compose:

    compose(
      compose(
        fs.createReadStream("./file.pcapng"),
        async function* (source) {
          for await (const chunk of source) {
            // Something here
            yield chunk
          }
        })
        .map((blockBuffer) => parseBufferBasedOnBlockType(blockBuffer))
        .filter((parsedBlock) => parsedblock.dest === "192.168.1.1"),
      async function* (source) {
        for await (const chunk of source) {
          // Something here too
          yield chunk
        }
      },
    )
      .toArray();

    And using it as a non-static method:

    fs.createReadStream("./file.pcapng")
      .compose(async function* (source) {
        for await (const chunk of source) {
          // Something here
          yield chunk
        }
      })
      .map((blockBuffer) => parseBufferBasedOnBlockType(blockBuffer))
      .filter((parsedBlock) => parsedblock.dest === "192.168.1.1")
      .compose(async function* (source) {
        for await (const chunk of source) {
          // Something here too
          yield chunk
        }
      })
      .toArray();
  5. mcollina commented on Jul 4, 2022

    @mcollina
    SponsorMember

    I think this use case is already covered by .map():

    fs.createReadStream("./file.pcapng")
      .map(async function (chunk) {
        // something async here
        return chunk
      })
      .map((blockBuffer) => parseBufferBasedOnBlockType(blockBuffer))
      .filter((parsedBlock) => parsedblock.dest === "192.168.1.1")
      .map(async function (chunk) {
        // something async here
        return chunk
      })
      .toArray();

    I'm not sure the syntax difference is worth adding another method on Readable. However, I'm happy to hear what other thinks.

  6. rluvaton commented on Jul 4, 2022

    @rluvaton
    MemberAuthor

    I think this use case is already covered by .map():

    fs.createReadStream("./file.pcapng")
      .map(async function (chunk) {
        // something async here
        return chunk
      })
      .map((blockBuffer) => parseBufferBasedOnBlockType(blockBuffer))
      .filter((parsedBlock) => parsedblock.dest === "192.168.1.1")
      .map(async function (chunk) {
        // something async here
        return chunk
      })
      .toArray();

    I'm not sure the syntax difference is worth adding another method on Readable. However, I'm happy to hear what other thinks.

    Map only map 1 input to 1 output and there are cases where we want to map 1/many input to many/1 output

  7. rluvaton commented on Jul 4, 2022

    @rluvaton
    MemberAuthor

    Let's say the first transformer is parsing binary file.

    Binary files are usually made out of blocks let's take for example the pcapng file format (Packet capture format that contains a "dump" of data packets captured over a network) as I'm very familiar with it:

    Pcapng file made out of blocks, each block has the total size in it as the size can change depending on the block body

                            1                   2                   3
        0 1 2 3 4 5 6 7 8 9 0 1 2 3 4 5 6 7 8 9 0 1 2 3 4 5 6 7 8 9 0 1
       +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+
     0 |                          Block Type                           |
       +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+
     4 |                      Block Total Length                       |
       +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+
     8 /                          Block Body                           /
       /              variable length, padded to 32 bits               /
       +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+
       |                      Block Total Length                       |
       +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+
    

    Using map we can't split the file to multiple blocks and parse each one individually.

  8. benjamingr commented on Jul 4, 2022

    @benjamingr
    Member

    @mcollina I think a good motivating example of something you can easily do with compose but not any of the other operators (or more accurately - I haven't found a way but maybe you know one) is readline.

    Consider this pseudocode for a readline implementation

    fs. createReadStream('./someFile').map(async function* lines(source) {
      let chunks = [];
      for await (const chunk of source) {
        const newLine = chunk.indexOf('\n')
        if(newLine !== -1) {
          yield Buffer.concat(chunksconcat(chunk.subarray(0, newLine);
          chunks = [chunk.subarray(newLine)];
       } else {
         chunks.push(chunk);
       }
      }
    }).forEach(processLine);

    The code can be made shorter but my point wasn't brevity but to underline how this fundamentally needs to keep state between iterations.

    For .map to be able to do this it'd need a "scope" it can reference between loop iterations easily e.g.:

    async function readlines(stream) {
      let chunks = []; // captured by closure
      return fs.createReadStream('./someFile').map(function (chunk) {
        const newLine = chunk.indexOf('\n')
        if(newLine !== -1) {
         yield Buffer.concat(chunksconcat(chunk.subarray(0, newLine);
         chunks = [chunk.subarray(newLine)];
       } else {
         chunks.push(chunk);
      }
      });
    }

    Now this technically works (unlike for example RxJS) because iterator helpers are on the iterator and not the iterable so the function is guaranteed to only run for one stream instance (which is actually very fortunate because this is a common footgun in .NET land) - but it's also less composable (the .map code is not reusable unlike the compose code) and relies on the function/block + a closure for scope.

  9. ronag commented on Jul 4, 2022

    @ronag
    Member

    .scan?

  10. rluvaton commented on Jul 4, 2022

    @rluvaton
    MemberAuthor

    I'm assuming you meant the RxJS scan method, I'm not very familiar with it but isn't it like reduce?

  11. ronag commented on Jul 6, 2022

    @ronag
    Member

    I'm assuming you meant the RxJS scan method, I'm not very familiar with it but isn't it like reduce?

    No, they are different.

    The upside of a compose is that one can have resource management inside the function, which is not possible with any of the other operators.

  12. ronag commented on Jul 6, 2022

    @ronag
    Member

    @rluvaton Would you mind opening a PR for .compose? Not promising we'll merge it, but I think it's worth further consideration.

  13. benjamingr commented on Jul 6, 2022

    @benjamingr
    Member

    .scan is like an iterative reduce with state. It's less powerful than .compose (name pending) like the static compose but I do think it can accommodate the use case mentioned here.

    That said I'd personally like to see more community interest before I'd feel we should do this?

  14. moved this to Pending Triage in Node.js feature requestson Oct 22, 2022
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

    feature requestIssues requesting new Node.js features.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