# node:stream

The `node:stream` module provides the Node.js abstract interface for streaming data: the `Readable`, `Writable`, `Duplex`, `Transform`, and `PassThrough` classes. A stream reads or writes data as a continuous flow of chunks, so a function processes a large volume of data without holding all of it in memory at once.

> **Note**
>
> Under `azion dev`, `import { pipeline } from "node:stream"` fails the build with `No matching export in "internal-env-dev:stream" for import "pipeline"`. The two examples that import `pipeline` run only in a deployed function.

---

## Examples

Each example is a complete function that imports from `node:stream`. The response below each example is the one a deployed function returns.

### Basic readable and writable streams

This function creates a readable stream and a writable stream, connects them with `pipe()`, and responds when the pipe is set up:

```javascript
/**
 * An example of using Node.js Stream API in an Azion Function.
 * Support:
 * - Partially supported (Extended by library `stream-browserify`)
 * @module runtime-apis/nodejs/stream/main
 * @example
 * // Build and run with the Azion CLI:
 * azion build
 * azion dev
 */
import stream from "node:stream";

/**
 * An example of using the Node.js Stream API in an Azion Function.
 * @param {*} event
 * @returns {Promise<Response>}
 */
const main = async (event) => {
  return new Promise((resolve, reject) => {
    const chunks = ["chunk1", "chunk2", "chunk3", "chunk4", "chunk5"];

    const nextChunk = () => {
      const chunk = chunks.shift();
      if (chunk) {
        console.log("Chunk", chunk);
        nextChunk();
      }
    };

    const readable = new stream.Readable({
      encoding: "utf8",
      read() {
        nextChunk();
      },
    });

    const writable = new stream.Writable({
      write(chunk, encoding, callback) {
        console.log("Chunk", chunk.toString());
        callback();
      },
    });

    readable.pipe(writable);

    resolve(new Response("Done"));
  });
};

export default main;
```

The function responds with:

```text
Done
```

### Transform stream for data processing

A `Transform` stream changes data as it passes through. This function connects two transforms and a collector with `pipeline()`: the first uppercases each line, the second numbers it:

```javascript
import { Transform, pipeline } from "node:stream";

const main = async (event) => {
  // Create a transform stream that uppercases text
  const upperCaseTransform = new Transform({
    transform(chunk, encoding, callback) {
      const uppercased = chunk.toString().toUpperCase();
      this.push(uppercased);
      callback();
    }
  });

  // Create a transform stream that adds line numbers
  let lineNumber = 0;
  const addLineNumbers = new Transform({
    transform(chunk, encoding, callback) {
      lineNumber++;
      const numbered = `${lineNumber}: ${chunk.toString()}`;
      this.push(numbered);
      callback();
    }
  });

  // Collect output
  const output = [];
  const collector = new Transform({
    transform(chunk, encoding, callback) {
      output.push(chunk.toString());
      callback();
    }
  });

  // Simulate input data
  const inputData = ["hello\n", "world\n", "azion\n", "runtime\n"];

  // Process each chunk through the pipeline
  for (const data of inputData) {
    upperCaseTransform.write(data);
  }
  upperCaseTransform.end();

  return new Promise((resolve) => {
    pipeline(
      upperCaseTransform,
      addLineNumbers,
      collector,
      (err) => {
        if (err) {
          console.error("Pipeline error:", err);
        }
        console.log("Output:", output.join(""));
        resolve(new Response(output.join("")));
      }
    );
  });
};

export default main;
```

The function responds with the numbered, uppercased lines:

```text
1: HELLO
2: WORLD
3: AZION
4: RUNTIME
```

### PassThrough stream for data forwarding

A `PassThrough` stream forwards data without changing it. This function writes three chunks to one and collects each chunk as it passes:

```javascript
import { PassThrough, pipeline } from "node:stream";

const main = async (event) => {
  // Create a PassThrough stream
  const passThrough = new PassThrough();

  // Collect data that passes through
  const collected = [];
  passThrough.on("data", (chunk) => {
    collected.push(chunk.toString());
    console.log("Data passed through:", chunk.toString());
  });

  // Write data to the stream
  const data = ["First chunk\n", "Second chunk\n", "Third chunk\n"];
  data.forEach(chunk => passThrough.write(chunk));
  passThrough.end();

  // Wait for all data to be processed
  await new Promise(resolve => passThrough.on("finish", resolve));

  return new Response(JSON.stringify({
    message: "Data passed through successfully",
    chunks: collected,
    totalChunks: collected.length
  }), {
    headers: { "Content-Type": "application/json" }
  });
};

export default main;
```

The function responds with the chunks it collected:

```json
{"message":"Data passed through successfully","chunks":["First chunk\n","Second chunk\n","Third chunk\n"],"totalChunks":3}
```

### Duplex stream for bidirectional communication

A `Duplex` stream is both readable and writable. This function writes two chunks to a duplex stream and reads three chunks from it:

```javascript
import { Duplex } from "node:stream";

const main = async (event) => {
  // Create a duplex stream
  const duplex = new Duplex({
    read(size) {
      // Simulate reading data
      if (this._readCounter === undefined) this._readCounter = 0;
      this._readCounter++;
      
      if (this._readCounter <= 3) {
        this.push(`Read chunk ${this._readCounter}\n`);
      } else {
        this.push(null); // End the stream
      }
    },
    write(chunk, encoding, callback) {
      console.log("Written:", chunk.toString().trim());
      callback();
    }
  });

  // Write to the duplex stream
  duplex.write("Data to write\n");
  duplex.write("More data\n");

  // Read from the duplex stream
  const output = [];
  duplex.on("data", (chunk) => {
    output.push(chunk.toString());
  });

  // Wait for stream to end
  await new Promise(resolve => duplex.on("end", resolve));

  return new Response(JSON.stringify({
    readOutput: output.join(""),
    message: "Duplex stream processed"
  }), {
    headers: { "Content-Type": "application/json" }
  });
};

export default main;
```

The function logs `Written: Data to write` and `Written: More data`, and responds with the chunks it read:

```json
{"readOutput":"Read chunk 1\nRead chunk 2\nRead chunk 3\n","message":"Duplex stream processed"}
```

### Request body processing with streams

This function reads the request body chunk by chunk, passes each chunk through two transforms, and counts the bytes:

```javascript
import { Transform } from "node:stream";
import { Buffer } from "node:buffer";

const main = async (event) => {
  const request = event.request;

  // Create a transform stream to process chunks
  let totalBytes = 0;
  const byteCounter = new Transform({
    transform(chunk, encoding, callback) {
      totalBytes += chunk.length;
      this.push(chunk);
      callback();
    }
  });

  // Create a transform to collect data
  const chunks = [];
  const collector = new Transform({
    transform(chunk, encoding, callback) {
      chunks.push(chunk);
      this.push(chunk);
      callback();
    }
  });

  // Get request body as stream (if available)
  if (request.body) {
    const reader = request.body.getReader();
    
    // Process the stream
    while (true) {
      const { done, value } = await reader.read();
      if (done) break;
      
      byteCounter.write(Buffer.from(value));
      collector.write(Buffer.from(value));
    }
    
    byteCounter.end();
    collector.end();
  }

  // Combine collected chunks
  const bodyContent = chunks.length > 0 
    ? Buffer.concat(chunks).toString("utf8") 
    : "No body content";

  return new Response(JSON.stringify({
    totalBytes,
    chunkCount: chunks.length,
    bodyPreview: bodyContent.substring(0, 200)
  }), {
    headers: { "Content-Type": "application/json" }
  });
};

export default main;
```

For a `POST` request whose body is `streamed request body for the stream sample`, the function responds with the byte count, the number of chunks, and the start of the body:

```json
{"totalBytes":43,"chunkCount":1,"bodyPreview":"streamed request body for the stream sample"}
```

---

## Supported APIs

The table lists the status of each `node:stream` API in Azion Runtime:

| API                    | Status                 |
| ---------------------- | ---------------------- |
| `stream.Readable`      | 🟢 Supported           |
| `stream.Writable`      | 🟢 Supported           |
| `stream.Duplex`        | 🟢 Supported           |
| `stream.Transform`     | 🟢 Supported           |
| `stream.PassThrough`   | 🟢 Supported           |
| `stream.pipeline()`    | 🟢 Supported           |
| `stream.compose()`     | 🟡 Partially supported |
| `stream.finished()`    | 🟡 Partially supported |
| `readable.pipe()`      | 🟢 Supported           |
| `readable.on('data')`  | 🟢 Supported           |
| `readable.on('end')`   | 🟢 Supported           |
| `readable.on('error')` | 🟢 Supported           |
| `writable.write()`     | 🟢 Supported           |
| `writable.end()`       | 🟢 Supported           |
| `transform.push()`     | 🟢 Supported           |

`stream.pipeline()` takes a callback, as the transform example shows. The promise form, `pipeline()` from `node:stream/promises`, throws `Error: [unenv] stream.promises.pipeline is not implemented yet!` in a deployed function and under `azion dev`.

---

## Related resources

- [Node.js APIs](/en/documentation/devtools/runtime/node.md): The status of every Node.js module in Azion Runtime, `stream` included.
- [Use Node.js APIs through polyfills](/en/documentation/guides/application-development/functions-and-runtime/use-polyfills.md): How the build turns `node:` imports into code that runs in Azion Runtime.
- [ReadableStream](/en/documentation/devtools/runtime/api-reference/readable-stream.md): The Web Streams alternative for reading data in chunks.
- [Node.js stream documentation](https://nodejs.org/api/stream.html): The full Node.js reference for every `node:stream` API in the table.
