Reference

std/stream/source

std/stream/src/source.trb

The reading end of a stream: Source, the stages built from it (map, take, through, ...), the terminal operations that read it to a value (toList, fold, into, ...), and its three ways to make one (from, pulling, produce).

trait Source

shared trait Source<Item, Failure> with Close

The reading end of a stream: the asynchronous sibling of Iterator, with the same verb and the same var self.

var lines = file.lines()
while const Some(line) = lines.next().await()? {
  print line
}

A source has an identity and is consumed once, so it is a shared type - and next changes it, so it takes var self exactly as Iterator.next does. Like every change of a shared object it goes through any binding that holds the source, a const one included ("Identity" in CONCEPT.md). That a var fn method may answer a Task is what the shared in shared type buys: for a value a var is an exclusive in-out access that ends with the call, for an object it is a permission to change the one object, and a permission survives an await.

  • A failure ends the stream. After Fail(problem) a source answers nothing useful any more; it may keep answering the same failure, and it never goes back to delivering items.
  • None ends the stream too, and it is final: after Ok(None) every further next() answers Ok(None).
  • One puller at a time. next() must not be called again before the task it answered has finished; two overlapping pulls are a bug, not a race the source has to defend against. A source may only be pulled by the task that made it - shared objects are confined to their task (CONCEPT, "Concurrency"), which is why nothing here spawns.
  • Backpressure is the pull. Nothing is read before somebody asks for it.
  • Whoever reads needs the permission - now or later. Reading is a var fn, and so does wrapping (map, filter, through, then, checked), because a wrapper keeps the source in a var field and pulls from it afterwards. So a const source cannot be consumed by wrapping it either, which is what a read-only view promised. A chain still reads as one expression, because a freshly produced object is a var path: nothing is thrown away where there is no copy, and nobody else holds a view of something that was just made.
  • close() releases what is above. The last release of a source runs it, and no program calls it: a reader that stops early (take, an abandoned loop, a failure it handles) lets go of the source, and every source built here holds the one it reads from in a field, which its own release releases after it. using is the form that pins the moment to the end of a block. close() is synchronous, cannot fail and may run after the stream ended; end() on the other end is the graceful counterpart that can.

Related

  • Sink - the writing end, and into/fill's other side.
  • Stage - what through wraps a source with.

fn next

var fn next(): Task<Result<Item?, Failure>>

Ok(Some(item)) for the next item, Ok(None) at the end, Fail(problem) for a failure that ends the stream.

fn through

var fn through<Output>(stage: Stage<Item, Output>): Source<Output, Failure>

The one place where a stage meets a stream. Everything below is this call with a name, and Json.items<User>() or lines() are the same stage values an Iterate takes.

fn map

var fn map<Output>(transform: Transform<Item, Output>): Source<Output, Failure>

Transforms every item as it is read.

fn filter

var fn filter(predicate: Predicate<Item>): Source<Item, Failure>

Reads only the items that answer the question.

fn filterMap

var fn filterMap<Output>(transform: Transform<Item, Output?>): Source<Output, Failure>

Transforms every item and drops the ones that answer None.

fn mapWhile

var fn mapWhile<Output>(transform: Transform<Item, Output?>): Source<Output, Failure>

Transforms items until one answers None, and ends the stream there.

fn take

var fn take(amount: Int): Source<Item, Failure>

Reads at most amount items, then ends the stream.

fn takeWhile

var fn takeWhile(predicate: Predicate<Item>): Source<Item, Failure>

Reads items while they answer the question, then ends the stream at the first one that does not.

fn skip

var fn skip(amount: Int): Source<Item, Failure>

Drops the first amount items, then reads the rest as they come.

fn indexed

var fn indexed(): Source<(index: Int, item: Item), Failure>

(0, first), (1, second), ..., with the halves named index and item as Iterate.indexed names them.

fn chunked

var fn chunked(size: Int): Source<List<Item>, Failure>

Groups items into lists of at most size; the last group is what is left.

fn then

var fn then<Output>(step: (value: Item) => Task<Result<Output, Failure>>): Source<Output, Failure>

The only asynchronous stage there is: a step that has to wait for each item (a lookup, a request, a write). A Stage cannot do this - it computes and never waits - so this is a method and not a stage value.

fn mapFailure

var fn mapFailure<Other>(transform: (failure: Failure) => Other): Source<Item, Other>

A stream of the same items with a failure of another type: what mapError is on a Result.

fn collect

var fn collect<Output>(into: Accumulator<Item, Output>): Task<Result<Output, Failure>>

The same Accumulators the pull side uses, so every terminal operation is written once. The values are pushed into a copy of it as they arrive, and isDone() ends the pull as soon as it has seen enough.

fn toList

var fn toList(): Task<Result<List<Item>, Failure>>

Every item, collected into a list.

fn count

var fn count(): Task<Result<Int, Failure>>

How many items the stream delivers.

fn fold

var fn fold<State>(initial: State, combine: (State, Item) => State): Task<Result<State, Failure>>

Combines every item into a running state, left to right, starting from initial.

fn forEach

var fn forEach(action: Action<Item>): Task<Result<Void, Failure>>

Runs action on every item as it arrives.

fn find

var fn find(predicate: Predicate<Item>): Task<Result<Item?, Failure>>

Stops at the first item that matches, which leaves the rest of the stream unread.

fn into

var fn into(var sink: Sink<Item, Failure>): Task<Result<Void, Failure>>

Hands everything to sink and finishes it. Both ends carry the same Failure on purpose: a pump that silently converted one into the other would hide which end broke. Use mapFailure on either side to say it.

fn from

static fn from(items: Iterate<Item>): Source<Item, Failure>

Everything an Iterate has, as a stream that never waits. Failure is whatever the caller needs: a source over a list cannot fail, and saying so as Never would make it unusable where a failing stream is expected.

fn pulling

static fn pulling(step: () => Task<Result<Item?, Failure>>): Source<Item, Failure>

The closure is next: the escape hatch for anything that already knows how to answer one item and has no state worth a type of its own. With state, write an ordinary shared type with var fields - that is what the traits here are for.

fn produce

static fn produce(capacity: Int = 0, body: (var sink: Sink<Item, Failure>) => Task<Result<Void, Failure>>): Source<Item, Failure> where Failure: From<ChannelClosed>

A producer that pushes instead of being pulled, without generators: body runs as a task of its own and writes into sink, and the items travel through a channel of capacity. The default 0 hands every item over directly, so the producer runs in lock-step with the consumer - which is what a generator does.

Source.produce { sink =>
  for row in query.rows() {
    sink.add(Report.of(row)).await()?
  }
  Ok void
}

Failure: From<ChannelClosed> because a consumer that stops pulling ends the producer, and that end needs a name in the producer's own failure type.

The body answers a Task, which is how it may await() the adds it writes: await() stands in a function whose result is a Task and in the closure spawn is handed, and nothing else in the language says that a callee runs a closure as a task. Its own body still produces the Result, exactly as the body of a fn that answers a Task does.

fn empty

static fn empty(): Source<Item, Failure>

A stream that ends immediately, without reading anything.

type Pulling

shared type Pulling<Item, Failure> with Source<Item, Failure>

Source.pulling { ... }: one closure for next, and the thing above it that has to be released.

fn of

static fn of(step: () => Task<Result<Item?, Failure>>, above: Close? = None): Pulling<Item, Failure>

Wraps step as a source, and above as what releasing it releases.

fn next

var fn next(): Task<Result<Item?, Failure>>

Calls step for the next item.

fn close

var fn close()

Nothing of its own: above is a field, and the release of the source releases it right after this.

type Staged

shared type Staged<Input, Item, Failure> with Source<Item, Failure>

source.through(stage). Two paths, like the synchronous Staged of std/iteration:

  • collect is fused: the stage wraps the collector's accumulator, no queue, one push per item.
  • next() needs a queue, because one item pushed in can become none or many coming out while the caller asks for exactly one. It holds what one item produced and is emptied before the next item is pulled.

The queue is a Shared<List<Item>> and not a plain field because an Accumulator is a value: the chain the stage built holds a copy of whatever it was given, so the only way to read what it pushed is a box both hold.

through on a staged source composes instead of stacking, so source.through(a).through(b) is one driver over a.then(b) and stays fused all the way down.

fn of

static fn of(var upstream: Source<Input, Failure>, stage: Stage<Input, Item>): Staged<Input, Item, Failure>

Wraps upstream with stage, and sets up the queue next() reads from.

fn next

var fn next(): Task<Result<Item?, Failure>>

Pulls from upstream and pushes into the stage until the queue holds an item or the stage is done.

fn through

var fn through<Output>(other: Stage<Item, Output>): Source<Output, Failure>

Composes other onto this stage instead of stacking a second Staged on top.

fn collect

var fn collect<Output>(into: Accumulator<Item, Output>): Task<Result<Output, Failure>>

Fused: no queue, and isDone() from the accumulator reaches the source through the stage.

fn close

var fn close()

Nothing of its own: upstream is a field, and the release of the source releases it right after this.

extend Source<Result<Item, Problem>, Failure>

extend<Item, Problem, Failure: From<Problem>> Source<Result<Item, Problem>, Failure>

A stage that can fail answers Result items (Json.items<User>(), lines()), because a Stage is synchronous and knows nothing about the stream's failure. checked() is where those items become the stream's failure and end it - the streaming counterpart of collecting an Iterate<Result<...>> into a Result<List<...>, ...>.

const all = body.through(Json.items<User>()).checked().toList().await()?

fn checked

var fn checked(): Source<Item, Failure>

Turns every item's own failure into the stream's failure, and ends the stream at the first one.