async stream experiment 4
by jib1
HTML
<div id="div"></div>
JavaScript
let base = performance.now();
const now = () => (performance.now() - base).toFixed();
const console = {log: msg => div.innerHTML += `${now()}: ${msg}<br>`};
const wait = ms => new Promise(resolve => setTimeout(resolve, ms));
const producerHighWaterMark = 3;
const processorHighWaterMark = 1;
(async url => {
try {
await new ChunkProducer()
.pipeThrough(new ChunkProcessor())
.pipeThrough(new ChunkProcessor())
.pipeTo(new WritableConsole());
} catch (e) {
console.log(`${e.name}: ${e.message}`);
}
})();
function ChunkProducer() {
return new ReadableStream({
chunks: [
{name: "A", close() { console.log(`closed ${this.name}`)}},
{name: "B", close() { console.log(`closed ${this.name}`)}},
{name: "C", close() { console.log(`closed ${this.name}`)}},
{name: "D", close() { console.log(`closed ${this.name}`)}},
{name: "E", close() { console.log(`closed ${this.name}`)}},
],
queued: [],
pull(controller) {
const chunk = this.chunks.shift();
const action = chunk? `Pulling ${chunk.name}` : "Flushing";
const {desiredSize} = controller;
console.log(`${action} (desiredSize=${desiredSize}, queue=${this.queued.length})`);
if (this.queued.length > producerHighWaterMark) this.queued.shift().close();
this.queued.push(chunk || {});
if (chunk) controller.enqueue(chunk);
},
}, {highWaterMark: producerHighWaterMark});
}
function ChunkProcessor() {
return new TransformStream({
async transform(chunk, controller) {
const name = chunk.name + chunk.name[0];
console.log(`Processing ${chunk.name} to ${name}...`);
await wait(1000);
console.log(`...done processing ${chunk.name} to ${name}`);
controller.enqueue({name});
}
}, {highWaterMark: processorHighWaterMark},
{highWaterMark: processorHighWaterMark});
}
function WritableConsole() {
return new WritableStream({
write: chunk => console.log(`Output: ${chunk.name}`),
...