Skip to content

Alternative end stream handling #2

Description

@robbiespeed

Having null chunks in transform adds extra branch handling to simple transforms, and shouldn't be needed in the stateful transform API.

Alternatively if chunks was never null, you could do end handling with the stateful API by simply adding end logic after consuming the source.

Altered examples from the article:

const toUpperCase = (chunks) => chunks.map(chunk => {
  const str = new TextDecoder().decode(chunk);
  return new TextEncoder().encode(str.toUpperCase());
});

// Stateful transform with resource cleanup
function createGzipCompressor() {
  // Hypothetical compression API...
  const deflate = new Deflater({ gzip: true });

  return {
    async *transform(source) {
      for await (const chunks of source) {
        for (const chunk of chunks) {
          deflate.push(chunk, false);
          if (deflate.result) yield [deflate.result];
        }
      }
      // Flush: finalize compression
      deflate.push(new Uint8Array(0), true);
      if (deflate.result) return [deflate.result];
    },
    abort(reason) {
      // Clean up compressor resources on error/cancellation
    }
  };
}

Similarly I'm not sure if abort is necessary, could it be a try/catch inside transform instead?

Activity

  1. jasnell commented on Feb 28, 2026

    @jasnell
    Collaborator

    Yep, considered this but also thinking about stateless transforms. I think it's important that the EOS signaling consistency is important. Not opposed to this but I'd like to see if we can have something that is consistent for both stateless and stateful transforms. Any thoughts for that?

  2. robbiespeed commented on Mar 1, 2026

    @robbiespeed
    Author

    There could be only a single type of transform (source: AsyncIterable<Uint8Array[]>) => AsyncIterable<TransformYield>. Any state or error handling can be accomplished inside the function itself:

    async function *gzipCompressor(source) {
      // Hypothetical compression API...
      // explicit resource management for automatic cleanup.
      using deflate = new Deflater({ gzip: true });
    
      try {
        for await (const chunks of source) {
          for (const chunk of chunks) {
            deflate.push(chunk, false);
            if (deflate.result) yield [deflate.result];
          }
        }
        // Flush: finalize compression
        deflate.push(new Uint8Array(0), true);
        if (deflate.result) return [deflate.result];
      } catch (reason) {
      }
    }

    Outside of making simple transforms more verbose, I'm not sure if this would block any optimizations. One way around this might be to introduce a seperate simpler transform API (chunk: Uint8Array) => Uint8Array:

    const uppercasedStream = Stream.mapChunks(source, (chunk) => {
      const str = new TextDecoder().decode(chunk);
      return new TextEncoder().encode(str.toUpperCase());
    });
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

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions