Skip to content

Bringing pull streams to modern JS #74

Description

@darsain

Pull streams are awesome, but their current implementation is limited to the JS from 10 years ago. This brings with it an awkward to use interface (callbacks everywhere), and limitations, some of which I described in #73.

This is a proposal to re-implement pull streams in modern JS.

Source streams

Source (readable) streams are functions that create an async generator.

Example:

values

// Transforms an iterable object (such as array) into an async generator.
async function * values(iterable) {
    for (const item of iterable) {
        yield item;
    }
}

Transform/through streams

Transform streams are functions that accept an async generator, and return another async generator.

Example streams:

map

// Applies `fn()` to all values before passing them through.
function map(fn) {
    return async function * (asyncGenerator) {
        for await (const value of asyncGenerator) {
            yield fn(value);
        }
    };
}

filter

// Lets through only values that pass the `test()`.
function filter(test) {
    return async function * (asyncGenerator) {
        for await (const value of asyncGenerator) {
            if (test(value)) {
                yield value;
            }
        }
    };
}

delay

// Delays each read by `ms`.
function delay(ms) {
    return async function * (asyncGenerator) {
        for await (const value of asyncGenerator) {
            await new Promise(resolve => setTimeout(resolve, ms));
            yield value;
        }
    };
}

errorOn

// Throws an error when value passing through matches `x`.
function errorOn(x) {
    return async function * (asyncGenerator) {
        for await (const value of asyncGenerator) {
            if (value === x) {
                throw new Error(`Value can't be "${x}"`);
            }
            yield value;
        }
    };
}

Sink (writable) streams

Sinks are functions that accept an async generator, and return a promise.

Example sinks:

log

// Logs all values pulled from passed stream.
async function log(asyncGenerator) {
    for await (const value of asyncGenerator) {
        console.log(value);
    }
}

reduce

// Reduces all values with reducer().
function reduce(reducer, acc) {
    return async function (asyncGenerator) {
        for await (const value of asyncGenerator) {
            acc = reducer(acc, value);
        }
        return acc;
    }
}

Examples

Assuming this little pull helper:

// Help us compose streams left to right.
function pull(...args) {
    let result = args.reverse().shift();
    for (let stream of args) {
        result = result(stream);
    }
    return result;
}

Read numbers 0 to 4 and sum them up:

const arr = Array(5).fill(0).map((x, i) => i); // array [0..4]
const numbers = values(arr); // source stream of numbers from arr
const sum = reduce((x, y) => x + y, 0); // sink to sum up all values

// raw
sum(numbers).then(x => console.log(x));

// with pull helper
pull(numbers, sum).then(x => console.log(x));

// in an async function
(async () => {
    console.log(await pull(numbers, sum));
})();

// All versions above log:
// > 10

Read numbers 0 to 4, delay each read by 300ms, filter only even numbers, and square the rest:

const arr = Array(5).fill(0).map((x, i) => i); // array [0..4]
const numbers = values(arr); // source stream of numbers from arr
const wait300 = delay(300); // transform stream to add 300 ms delay to each read
const even = filter(x => x % 2 === 0); // transform stream to filter only even numbers
const square = map(x => x * x); // transform stream to square each number

// raw
log(square(even(wait300(numbers)))).then(() => console.log('done'));

// with pull helper
pull(numbers, wait300, even, square, log).then(() => console.log('done'));

// in an async function
(async () => {
    await pull(numbers, wait300, even, square, log);
    console.log('done');
})();

// All versions above log:
// > 0
// > 4
// > 16
// > done

Read numbers 0 to 4, but throw an error on 2:

const arr = Array(5).fill(0).map((x, i) => i);

pull(values(arr), errorOn(2), log).then(null, err => console.log(`Error: ${err.message}`));

// in an async function
(async () => {
    try {
        await pull(values(arr), errorOn(3), log);
    } catch(err) {
        console.error(`Error: ${err.message}`);
    }
})();

// All versions above log:
// > 0
// > 1
// > Error: Value can't be "2"

Comparison with current implementation

The implementation above has it all. Backpressure, error propagation, and more:

Simpler to work with

You don't need to be passing callbacks into sinks and doing other stuff from 10 years ago.

Do you want the final value from reducer? Just await, or .then() it to something else.

Do you want to catch errors? try/catch or add .then(null, onError) handler.

It all just seamlessly fits into modern async control flow.

Easily identifiable interfaces

Sources and transforms are async generators, while sinks are promises. We can duck type and work with all of these interfaces already. No flags required.

Simpler to implement

Everything looks like sync code. No callbacks juggling whatsoever. Non-blocking by default due to nature of promises. No fear of blowing some stinky call stack. Just look how simple each stream constructor looks like!

Easier to understand and reason about

Once you understand async functions & async generators, the brain power required to thinking about, use, and implement these async interfaces drops to an absolute minimum. No callbacks in sight.

Environment support

async/await is gonna land in node behind a flag October 2016 (this month), and without a flag possibly mid December (already in Chrome Canary).

Async generators is currently stage 3 (candidate), which means the spec is complete, and browser implementations are going to be landing any moment.

Regardless, you can use all of it already with babel transpilation. Just try and play with any of the code above in babel repl.


And not a single callback was passed that day.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions