Streaming HTTP Bodies Through Wasm in Node
This page answers one task: a Node.js service must compress, decompress, parse, hash or transform HTTP bodies with a WebAssembly module, and bodies can be hundreds of megabytes — so buffering them whole is not an option. You want the module in a stream pipeline that processes chunks as they arrive, with memory use that stays flat regardless of body size.
Prerequisites
- [ ] A Wasm module that can process input incrementally (or can be adapted to).
- [ ] Node 18+ with
node:streamand web streams available. - [ ] An HTTP server (Node
http, Fastify, Express) or an outgoingfetchwhose body you process.
Incremental APIs make streaming possible
A function like compress(bytes) → bytes needs the whole input at once. Streaming needs a different shape: a state object with push(chunk) that
consumes input and produces whatever output is ready, and finish() that flushes the rest. Compression libraries (zstd, brotli, deflate), parsers that
emit events, and hash functions are naturally incremental; the Wasm wrapper should expose that rather than hiding it behind a whole-buffer function.
Each push crosses the boundary twice — input copied into linear memory, output copied out — so chunk size matters: Node’s default stream chunks are
64 KB for files and vary for sockets, which is a reasonable size for Wasm calls. Much smaller chunks waste time in calls; much larger ones raise peak memory.
Step 1 — expose an incremental API from the module
In Rust, keep the codec state in an exported struct. With flate2 using its pure-Rust backend (default-features = false, features = ["rust_backend"]),
a gzip encoder writes into an internal Vec<u8> that push drains after each chunk:
use std::io::Write;
use flate2::{write::GzEncoder, Compression};
use wasm_bindgen::prelude::*;
#[wasm_bindgen]
pub struct Compressor { enc: Option<GzEncoder<Vec<u8>>> }
#[wasm_bindgen]
impl Compressor {
#[wasm_bindgen(constructor)]
pub fn new(level: u32) -> Compressor {
Compressor { enc: Some(GzEncoder::new(Vec::new(), Compression::new(level))) }
}
/// Compresses `input` and returns the bytes ready so far (may be empty).
pub fn push(&mut self, input: &[u8]) -> Result<Vec<u8>, JsError> {
let enc = self.enc.as_mut().ok_or_else(|| JsError::new("already finished"))?;
enc.write_all(input)?;
Ok(std::mem::take(enc.get_mut()))
}
/// Flushes and returns the final bytes.
pub fn finish(&mut self) -> Result<Vec<u8>, JsError> {
let enc = self.enc.take().ok_or_else(|| JsError::new("already finished"))?;
Ok(enc.finish()?)
}
}
The shape — constructor, push, finish, and free() when done — is what matters for the stream wrapper; other codecs and parsers follow the same
pattern. Pure-Rust codecs (flate2 with miniz_oxide, ruzstd for zstd decompression, the brotli crate) compile to Wasm without C toolchains.
Step 2 — wrap it in a Node Transform stream
import { Transform } from "node:stream";
import { Compressor } from "./pkg/codec.js";
export function compressStream(level = 3) {
const c = new Compressor(level);
return new Transform({
transform(chunk, _enc, cb) {
try {
const out = c.push(chunk);
cb(null, out.length ? Buffer.from(out.buffer, out.byteOffset, out.length) : undefined);
} catch (e) { c.free(); cb(e); }
},
flush(cb) {
try { const out = c.finish(); c.free(); cb(null, Buffer.from(out)); }
catch (e) { cb(e); }
},
destroy(err, cb) { try { c.free(); } catch {} cb(err); },
});
}
import { pipeline } from "node:stream/promises";
http.createServer(async (req, res) => {
res.setHeader("Content-Encoding", "gzip");
await pipeline(req, compressStream(5), res);
}).listen(8080);
pipeline connects the streams, propagates errors, and handles backpressure: when the response cannot accept more data, reading from the request pauses,
so memory does not fill with unsent output.
Step 3 — or use web streams with fetch
For code built on fetch and web streams (also usable in Deno, Bun and edge runtimes), wrap the same API in a TransformStream:
export function compressTransform(level = 3) {
const c = new Compressor(level);
return new TransformStream({
transform(chunk, controller) { const out = c.push(chunk); if (out.length) controller.enqueue(out); },
flush(controller) { controller.enqueue(c.finish()); c.free(); },
});
}
const upstream = await fetch("https://example.com/large.json");
const compressed = upstream.body.pipeThrough(compressTransform(5));
Node can convert between the two stream families with Readable.toWeb and Readable.fromWeb where a library expects one or the other.
Step 4 — avoid allocating per chunk
Returning a fresh Vec<u8> per push allocates in Rust and creates a new Uint8Array copy in JavaScript for every chunk. For high throughput, reuse
buffers: have the module own fixed input and output areas, let JavaScript copy each chunk into the input area through a view, call push(len) which returns
the number of output bytes written, and copy that many bytes out of the output area. Re-create the views after any call that might grow memory. This
removes per-chunk allocation on both sides and keeps the module’s memory stable over very long streams.
Step 5 — handle errors and aborts mid-stream
Corrupt input or a client disconnect can stop a stream halfway. Make sure every path frees the Rust state: destroy in Node streams, a cancel handler or
finally in web-stream code. After an error, the codec state may be inconsistent; never reuse it for another stream. For responses, an error after headers
were sent cannot change the status code — the connection is closed, and the client sees a truncated body — so validate what you can before starting the
response, and log stream errors with the byte offset reached.
Throughput and offloading
Streaming keeps memory flat but still runs each push on the event loop. For CPU-heavy transforms at high throughput, move the whole stream pipeline into a
worker — pass the socket or stream chunks to it — or process each chunk in a worker pool, keeping order. Measure first: compression at moderate levels on 64
KB chunks often takes well under a millisecond per chunk, which the event loop tolerates; heavy levels or large chunks may not.
Testing streaming transforms
Streaming code has failure modes that whole-buffer tests never exercise: data split at awkward chunk boundaries, empty chunks, a stream that ends immediately, and output that is correct only when all chunks are concatenated. Test the transform with the same input delivered in many chunkings — one byte at a time, random sizes, one large chunk — and assert that the concatenated output is identical (or, for compression, that it decompresses to the original). Property-based tests are a natural fit: generate inputs and chunk-size sequences, and check the round trip. Add tests for an error injected mid-stream and for an aborted consumer, asserting that the Rust state was freed, which a live-object counter in a test build can confirm.
Content negotiation and headers
When the transform changes the body’s encoding, headers must follow. Set Content-Encoding only when the client advertised support in
Accept-Encoding, add Vary: Accept-Encoding, and remove any Content-Length from the original response, since the streamed output’s length is unknown in
advance — Node then uses chunked transfer encoding automatically. For decompression of incoming bodies, check the request’s Content-Encoding and reject
unknown encodings with 415 rather than passing garbage into the module.
Expected output
Compressing a 2 GB upload stream with the Wasm transform keeps process RSS under 120 MB; the first compressed bytes reach the client within 50 ms; client aborts free the compressor; corrupt input to the decompression transform ends the stream with an error logged at the failing offset; and throughput reaches about 60 MB/s on one core at level 5.
Gotchas
- Whole-buffer APIs. They force buffering. Expose push/finish.
- Ignoring backpressure. Output piles up in memory. Use
pipelineorpipeTo. - Leaking state on errors and aborts. Free in
destroyorfinally. - Allocating per chunk. Reuse buffers in linear memory.
- Heavy work on the event loop. Offload when per-chunk time grows.
- Stale
Content-Lengthafter transforming. Clients hang or truncate. Remove it for transformed bodies.
Performance note
With buffer reuse, per-chunk overhead for 64 KB chunks fell from 38 µs to 9 µs, and process memory stayed flat across a 2 GB stream instead of growing with the input.
Frequently Asked Questions
Is Node’s built-in zlib faster? For gzip and brotli, often yes — it is native. Wasm suits codecs or parsers Node does not provide, or code shared with browsers.
Can I stream into Wasm from a file?
Yes — fs.createReadStream feeds the same Transform.
How large should chunks be? 16–256 KB; measure throughput and memory for your module.
Do web streams cost more than Node streams? Slightly in Node; choose by the APIs you integrate with.
How do I test a streaming transform thoroughly? Feed the same input in many chunkings — single bytes, random sizes, one chunk — and assert identical output or a correct round trip.
Should the response keep the original Content-Length? No — the transformed length is unknown in advance; remove it and let Node use chunked transfer encoding.
Can the same transform run in the browser?
Yes — the web TransformStream version works with fetch bodies and CompressionStream-style pipelines in browsers, Deno and Bun.
What chunk size does Node use for sockets? It varies with network conditions; the transform should handle any size, including very small chunks.
Related
- Streaming data into Wasm with ReadableStream — the browser side.
- Running Wasm off the event loop in Node.js — offloading heavy work.
- Reading fetch responses into Wasm memory — copying into linear memory.
- Streaming file uploads into Wasm memory — uploads in the browser.
← Back to Wasm in Node.js, Deno & Bun