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.