Streaming Data into Wasm with ReadableStream

This guide answers one task: process a large input — an upload, a download, a generated sequence — by feeding it into a WebAssembly module in chunks, so peak memory stays bounded and the first result appears before the last byte arrives.

Prerequisites

  • [ ] A module that can consume input incrementally, or one you can change to.
  • [ ] A ReadableStream source: fetch’s body, File.stream(), or your own.
  • [ ] A fixed input region in linear memory, sized once.
  • [ ] A reason: the input is large enough that buffering it matters.

Why not just buffer it

The straightforward version reads the whole input into an ArrayBuffer, copies it into the module, and calls once. For a few megabytes that is correct and simpler, and it should be the default.

It stops working at scale for two reasons. Peak memory becomes the input plus the module’s copy plus whatever the module allocates — comfortably three times the file — which a phone will refuse for a large upload. And nothing happens until the last byte arrives, so a forty-second download is forty seconds of nothing followed by a result.

Streaming fixes both, at the cost of an interface that carries state between calls.

Peak memory, and when the first result appears Buffering holds the whole input on both sides before any work begins. Streaming keeps one chunk resident and emits results as they become available, so memory is bounded and output starts early. buffered download whole input copy into module process first result streamed chunk chunk chunk first result here, and results continue as chunks arrive Peak memory is one chunk plus the module's own state, whether the input is ten megabytes or ten gigabytes.

The module’s side

A streaming module needs three entry points and a fixed region the host writes into.

const CHUNK: usize = 1 << 20;                       // 1 MiB, decided once
static mut INPUT: [u8; CHUNK] = [0; CHUNK];
static mut PARSER: Option<Parser> = None;

#[no_mangle] pub extern "C" fn input_ptr() -> *mut u8 { unsafe { INPUT.as_mut_ptr() } }
#[no_mangle] pub extern "C" fn input_cap() -> usize { CHUNK }

#[no_mangle]
pub extern "C" fn begin() -> i32 { unsafe { PARSER = Some(Parser::new()) }; 0 }

/// Consume `len` bytes from the input region. Returns the number of output bytes ready, or a negative code.
#[no_mangle]
pub extern "C" fn feed(len: usize) -> i32 {
    let p = unsafe { PARSER.as_mut().unwrap() };
    match p.consume(unsafe { &INPUT[..len] }) {
        Ok(out_len) => out_len as i32,
        Err(e) => -(e.code() as i32),
    }
}

#[no_mangle]
pub extern "C" fn finish() -> i32 { /* flush any buffered tail */ 0 }

A fixed array rather than a heap allocation matters: it lives at a stable address in the data segment, so the host’s view over it never needs rebuilding and the module never grows memory mid-stream.

The parser keeps whatever partial state a chunk boundary leaves — half a record, an incomplete multi-byte sequence — which is the part that makes this design different from a single-call one.

The host’s side

Read from the stream, copy each chunk into the region, call feed, and drain whatever output appeared.

export async function streamThrough(mod, stream, onOutput) {
  const ptr = mod.exports.input_ptr();
  const cap = mod.exports.input_cap();
  const view = new Uint8Array(mod.exports.memory.buffer, ptr, cap);

  mod.exports.begin();
  const reader = stream.getReader();
  try {
    for (;;) {
      const { value, done } = await reader.read();
      if (done) break;
      for (let off = 0; off < value.length; off += cap) {
        const slice = value.subarray(off, Math.min(off + cap, value.length));
        view.set(slice);
        const n = mod.exports.feed(slice.length);
        if (n < 0) throw new Error(`module error ${n}`);
        if (n > 0) onOutput(readOutput(mod, n));
      }
    }
    const tail = mod.exports.finish();
    if (tail > 0) onOutput(readOutput(mod, tail));
  } finally {
    reader.releaseLock();
  }
}

The inner loop exists because a stream chunk is whatever size the source produced — often 64 kB from a network, sometimes megabytes from a file — and the module’s region has a fixed capacity. Splitting the chunk rather than resizing the region keeps memory flat.

Awaiting read() is what yields to the event loop, so this loop is naturally cooperative: the page renders between chunks without any additional scheduling.

Chunk boundaries, which the module must handle

A chunk almost never ends at a record boundary, and the module has to cope. Three patterns cover most formats.

Carry the tail. The parser keeps unconsumed bytes and prepends them to the next chunk. Simple, correct, and the carry buffer must be bounded — a format that permits an unbounded record means an attacker can make the carry grow without limit.

fn consume(&mut self, chunk: &[u8]) -> Result<usize, Error> {
    self.carry.extend_from_slice(chunk);
    if self.carry.len() > MAX_RECORD { return Err(Error::RecordTooLarge); }
    let mut consumed = 0;
    while let Some((rec, n)) = try_parse(&self.carry[consumed..]) {
        self.emit(rec); consumed += n;
    }
    self.carry.drain(..consumed);
    Ok(self.output_len())
}

Ask for a specific length. For a length-prefixed format the parser knows exactly how many bytes it needs next, and can tell the host, which then supplies precisely that. More round trips, no carry buffer.

Require aligned chunks. Where the host controls the source, feeding whole records is simplest of all — and is why an interface that accepts arbitrary chunks is strictly more useful than one that does not.

The tail that crosses a boundary A record beginning near the end of one chunk continues in the next. The parser keeps the partial bytes in a bounded carry buffer and completes the record once the following chunk arrives. chunk 1 record A record B C (part) carry C (part) held until the next chunk, bounded by MAX_RECORD chunk 2 C (rest) record D C completes and is emitted Without a bound on the carry, an input claiming an enormous record makes the buffer grow until allocation fails — which is a denial of service, not a parse error.

Choosing which pattern fits

The carry approach suits self-delimiting formats — newline-separated records, tag-length-value frames, anything where the parser can recognise a complete unit by looking at the bytes. It is the most forgiving of the three because the host needs to know nothing about the format at all.

Asking for a specific length suits formats with an explicit header that states the body size. The module becomes a small state machine — want 4 bytes, now want N bytes, repeat — and never buffers more than one record. The cost is that the host loop grows a second dimension: it must satisfy a request that may span several stream chunks, or be satisfied several times over from one.

Requiring aligned chunks is only safe where the same code produced the input. It is common inside a single application — a worker generating frames for a module in the same page — and unsafe the moment the input comes from a network or a file the user chose.

Backpressure and output

If the module produces output faster than the consumer accepts it, something has to give. A TransformStream expresses that naturally and gives the consumer control.

export function wasmTransform(mod) {
  return new TransformStream({
    start() { mod.exports.begin(); },
    async transform(chunk, controller) {
      const view = inputView(mod);
      for (let off = 0; off < chunk.length; off += view.length) {
        view.set(chunk.subarray(off, off + view.length));
        const n = mod.exports.feed(Math.min(view.length, chunk.length - off));
        if (n > 0) controller.enqueue(readOutput(mod, n));
      }
    },
    flush(controller) {
      const tail = mod.exports.finish();
      if (tail > 0) controller.enqueue(readOutput(mod, tail));
    },
  });
}
await file.stream().pipeThrough(wasmTransform(mod)).pipeTo(destination);

The stream machinery then handles backpressure: if the destination is slow, transform is not called again until it catches up, and the source is not read. That is considerably better behaviour than a manual loop that reads as fast as it can and buffers the difference.

Expected output

Streaming shows a flat memory profile and output starting early:

input: 2.1 GB
chunk: 1 MiB
  chunk    1/2148   heap 14 MB   first output at 41 ms
  chunk 1000/2148   heap 14 MB
  chunk 2148/2148   heap 15 MB
done in 38.2 s, peak heap 15 MB

Against the buffered version on the same input, which did not complete: RangeError: Array buffer allocation failed at roughly 2 GB. The flat heap column is the property to verify — a number that climbs with chunk index means something is being retained per chunk.

What the host still owes the module

Two responsibilities do not move into the stream machinery. The first is cancellation: if the user navigates away or aborts, the reader must be cancelled and the module reset, or the next run starts on top of half-consumed state. Wiring an AbortSignal to reader.cancel() and calling begin() afresh on the next run covers it.

The second is error propagation. A negative return from feed means the input was malformed, and the right response is usually to stop rather than to skip the chunk and continue — a parser that has lost its place produces garbage for the rest of the stream. Throwing from transform errors the stream, which propagates to the destination and to whatever is awaiting pipeTo, and that is the behaviour you want.

The module as a transform Wrapping the module in a TransformStream hands backpressure to the stream machinery: a slow destination stops the source being read, with no manual buffering. source stream fetch or File transform feed + drain module fixed input region destination writes or renders Backpressure flows right to left: when the destination is slow, transform is not called and nothing is read. Peak memory is one chunk plus the module's own state, whatever the size of the input. Cancellation has to be wired explicitly — an aborted read must also reset the module's parser state.

Gotchas

  • Unbounded carry buffer. A crafted input grows it until allocation fails.
  • Assuming stream chunks are a useful size. They are whatever the source produced; split them yourself.
  • Rebuilding the input view per chunk. Unnecessary if the module never grows memory, and a correctness requirement if it does.
  • Reading output only at the end. Defeats the purpose; drain after every feed.
  • Ignoring the reader’s lock. Release it in a finally, or a later reader cannot attach.
  • No finish call. Buffered tail data is silently dropped.

Performance note

Streaming a 2.1 GB file through a 1 MiB region held peak heap at 15 MB and sustained about 55 MB/s, of which the boundary crossings accounted for under 2% — 2,148 calls in total. The buffered equivalent could not run at all above roughly 1.5 GB. Raising the chunk to 4 MiB improved throughput by about 4% and raised peak memory by 3 MB, which is a reasonable trade and not a dramatic one.

Frequently Asked Questions

What chunk size should I use? One to four mebibytes suits most workloads. Smaller increases per-chunk overhead, larger increases peak memory for diminishing throughput. Measure once on a realistic input.

Does this work in a worker? Yes, and it is the better place for it. Streams are available in workers, and the await in the loop yields there just as it does on the main thread — while the actual processing blocks nobody.

Can the module pull rather than be pushed? With a suspending import it can call back for more data, as described in calling async JavaScript with JSPI. The push design here needs no proposal and works everywhere, which is why it is the default.

← Back to Async & Event-Loop Integration