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
ReadableStreamsource:fetch’sbody,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.
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.
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.
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
finishcall. 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.
Related
- Working with datasets larger than memory — the same discipline for computation rather than input.
- Keeping the UI responsive during long Wasm tasks — chunking the work as well as the data.
- Reading Wasm linear memory with typed arrays — the view mechanics this relies on.
← Back to Async & Event-Loop Integration