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

Writing your own steps

When no method does what you need, write the step as an async generator function and pass it to transform(). It gets the whole sequence, so it can keep state between items and yield as many or as few as it likes.

import { enumerate, read, type TransformerFunction } from "@j50n/proc";

// Fewer items, with state: drop a line that repeats the one before it.
async function* dedupe(lines: AsyncIterable<string>) {
  let previous: string | undefined;
  for await (const line of lines) {
    if (line !== previous) yield line;
    previous = line;
  }
}

// Batching: a function that makes the step, so the size can vary.
function batches<T>(size: number): TransformerFunction<T, T[]> {
  return async function* (items) {
    let batch: T[] = [];
    for await (const item of items) {
      batch.push(item);
      if (batch.length === size) {
        yield batch;
        batch = [];
      }
    }
    if (batch.length > 0) yield batch; // the last, short batch
  };
}

console.log(
  await enumerate(["a", "a", "b", "a", "a"]).transform(dedupe).collect(),
);
console.log(await read("fruit.txt").lines.transform(batches(2)).collect());
[ "a", "b", "a" ]
[ [ "cherry", "apple" ], [ "banana" ] ]

A step is any function from an AsyncIterable<T> to an AsyncIterable<U>; the type TransformerFunction<T, U> names it. A generator function is the step itself: pass dedupe, not dedupe(). For a step that takes settings, write a function that returns the step, and call that: batches(2). The built-in steps come in both kinds (toBytes and gzip are passed as they are; buffer(size) and the data formats’ fromCsvToRows() are called); the list is in Enumerable and its methods.

More items than came in

yield* hands on every item of an iterable:

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

// More items: each line becomes its words.
async function* words(lines: AsyncIterable<string>) {
  for await (const line of lines) {
    yield* line.split(/\s+/).filter((word) => word.length > 0);
  }
}

console.log(
  await enumerate(["the quick", "  brown fox "]).transform(words).collect(),
);
[ "the", "quick", "brown", "fox" ]

For a step that only changes or drops single items, map, filter, and flatMap are shorter. Reach for a generator when the step needs memory of earlier items, needs to see the end (to flush a last batch, or emit a total), or needs to stop early.

Stopping early

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

// A step that returns before its input ends closes the source behind it.
async function* untilEnd(lines: AsyncIterable<string>) {
  for await (const line of lines) {
    if (line === "END") return;
    yield line;
  }
}

const lines = await run("sh", "-c", "echo a; echo b; echo END; echo c; exit 3")
  .lines
  .transform(untilEnd)
  .collect();
console.log(lines); // no error: the exit code isn't checked after an early stop
[ "a", "b" ]

Returning from the generator before its input ends closes the source, as take does: the command’s output is closed, and its exit code isn’t checked. Code in a finally block of the generator runs when the step stops for any reason, including a consumer that stops early downstream.

Errors

An error thrown in a step comes out of the consumer’s await, unchanged. An error from upstream (a failed command, a callback that threw) is thrown inside the step, from its for await loop, so a step can catch it, to add context or to recover:

import { ExitCodeError, run } from "@j50n/proc";

// A step sees errors from upstream in its own loop, and can add context.
async function* numbered(lines: AsyncIterable<string>) {
  let n = 0;
  try {
    for await (const line of lines) {
      n += 1;
      yield `${n}: ${line}`;
    }
  } catch (error) {
    throw new Error(`input failed after line ${n}`, { cause: error });
  }
}

try {
  await run("sh", "-c", "echo a; echo b; exit 2")
    .lines
    .transform(numbered)
    .forEach((line) => console.log(line));
} catch (error) {
  if (error instanceof Error) {
    console.log(error.message);
    if (error.cause instanceof ExitCodeError) {
      console.log(`cause: exit code ${error.cause.code}`);
    }
  }
}
1: a
2: b
input failed after line 2
cause: exit code 2

Rethrow with the original as cause rather than swallowing it. A step that catches an error and simply returns hides the failure: the pipeline ends as if the input had run out.

With a TransformStream

transform() also takes a TransformStream, or any { writable, readable } pair, such as CompressionStream:

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

// A TransformStream works too.
function upper() {
  return new TransformStream<string, string>({
    transform(line, controller) {
      controller.enqueue(line.toUpperCase());
    },
  });
}

console.log(await enumerate(["one", "two"]).transform(upper()).collect());

// A stream is used up after one pass: the second use yields nothing.
const once = upper();
console.log(await enumerate(["three"]).transform(once).collect());
console.log(await enumerate(["four"]).transform(once).collect());
[ "ONE", "TWO" ]
[ "THREE" ]
[]

A stream works once. Used a second time, it yields nothing or throws, and the source goes unread, so create a new one for each use, as upper() does. An error from upstream of the stream reaches the consumer unchanged, and so does one thrown in its transform(). A generator is usually simpler to write; use a stream when you already have one.