Keyboard shortcuts

Press ← or → to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

From callbacks to iterables

WritableIterable is a queue: push-style code (callbacks, event handlers, timers) writes items into it, and a for await loop or an Enumerable chain reads them out.

import { enumerate, WritableIterable } from "@j50n/proc";

// A producer that calls back on a timer, as many event sources do.
const ticks = new WritableIterable<number>();
let n = 0;
const timer = setInterval(() => {
  n += 1;
  ticks.write(n);
  if (n === 3) {
    clearInterval(timer);
    ticks.close(); // without this, the loop below waits forever
  }
}, 10);

await enumerate(ticks)
  .map((tick) => `tick ${tick}`)
  .forEach((line) => console.log(line));
console.log("closed");
tick 1
tick 2
tick 3
closed

It has three operations:

  • write(item) adds an item to the queue and returns at once.
  • close() ends the data. The reader gets every item written, then its loop ends.
  • close(error) ends the data with an error. The reader gets every item written before it, then the error is thrown from its loop.

Only the first close counts; later ones are ignored. isClosed tells you whether it has happened.

Bridging events

import { WritableIterable } from "@j50n/proc";

const target = new EventTarget();
const messages = new WritableIterable<string>();

target.addEventListener("message", (event) => {
  // write() rejects once the queue is closed; catch it in a handler, or the
  // unhandled rejection ends the program.
  messages.write((event as CustomEvent<string>).detail).catch(() => {});
});
target.addEventListener("error", () => {
  messages.close(new Error("the source failed"));
});

// Nothing is reading yet, so these just queue.
target.dispatchEvent(new CustomEvent("message", { detail: "hello" }));
target.dispatchEvent(new CustomEvent("message", { detail: "world" }));
target.dispatchEvent(new Event("error"));
target.dispatchEvent(new CustomEvent("message", { detail: "too late" }));

try {
  for await (const message of messages) {
    console.log(message);
  }
} catch (error) {
  if (error instanceof Error) console.log(`error: ${error.message}`);
}
hello
world
error: the source failed

The handlers run before anything reads, so the items wait in the queue until the loop starts. The error event closes the queue with an error, which the reader sees after “hello” and “world”; the message after that is refused.

Feeding a command

A WritableIterable can be a command’s stdin. .run() starts the command at once and passes items to it as they are written:

import { enumerate, WritableIterable } from "@j50n/proc";

// The command starts now, and reads its stdin as items are written.
const input = new WritableIterable<string>();
const sorted = enumerate(input).run("sort").lines.collect();

for (const fruit of ["pear", "apple", "fig"]) {
  await input.write(fruit);
}
await input.close(); // sort sees the end of its input

console.log(await sorted);
[ "apple", "fig", "pear" ]

Enumerable.writeTo() also accepts a WritableIterable, to copy a sequence into one.

Traps

Nothing closes it for you. Until close() is called, the reader waits for the next item. If something else keeps the program alive (a timer, a server), it waits forever; if nothing does, Deno stops with error: Top-level await promise never resolved. Make sure every way the source can end, including failure, leads to a close.

There is no backpressure. write() doesn’t wait for the reader; its promise resolves right away. When the writer is faster than the reader, or nothing reads at all, the queue grows without limit and every item stays in memory. That suits events, which can’t be paused anyway. If the producer can wait, and the data is large, make it a generator instead and let the reader pull.

write() after close() rejects. In an event handler, nobody awaits that promise, so the rejection is unhandled and ends the program. Catch it, as the example does, or check isClosed first.

A reader that stops doesn’t stop the producer. After a break, take, or first, write() still succeeds, but the items go nowhere. If the producer should stop when the reader does, close the queue yourself when the loop ends, in a finally, and have the producer check isClosed.

One reader, once. Read it with a single loop or chain. A second reader, at the same time or after the first, throws TypeError: a WritableIterable can be read only once.

The constructor takes an onclose callback, called on the first close(), which waits for it: a place to remove event listeners or clear a timer. See the API reference.