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. Noneends the stream too, and it is final: afterOk(None)every furthernext()answersOk(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 avarfield and pulls from it afterwards. So aconstsource 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 avarpath: 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.usingis 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
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:
collectis 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.