std/stream/sink
std/stream/src/sink.trb
The writing end of a stream: the asynchronous sibling of Accumulator, with the same two verbs and the same
var self.
trait Sink
shared trait Sink<Item, Failure> with Close
using file = File.create(path)?
file.addLine("first").await()?
file.end().await()
A sink has an identity and is ended once, so it is a shared type, and add and end change it: they take
var self, exactly as Accumulator.add does. A sink that is written to sits in a var binding, and a const
handle is the read-only view every shared object has. A var fn method may answer a Task because for an object
var is a permission and not an exclusive access that ends with the call (CONCEPT, "Identity").
- Backpressure is the
awaitonadd. The task finishes when the target has taken the item, not before, so a writer that is faster than its target waits by itself and nothing piles up. - Order is the order of the calls.
addmust not be called again before the task it answered has finished. end()is the graceful end: everything buffered is written, the target learns that no more is coming, and whatever went wrong on the way is reported here at the latest. Afterend(),addfails. The word is notfinish, which belongs toAccumulatorand answers a result: the end of a stream is not a value.- A failure ends the stream. After a failure neither
addnorendwill succeed again. close()is the abrupt end: it releases the target and cannot fail, which is whatCloseneeds: the last release of the sink runs it, and no program calls it. A sink that is released withoutend()may have written less than it was given - that is the point of the two.- Buffering is never implicit.
sink.buffered(capacity:)says it, andBuffered.flush()is where it is forced. - Whoever writes needs the permission - now or later.
mapFailureandbufferedtakevar selfas well, because the wrapper keeps the sink in avarfield and writes into it afterwards. A chain is still one expression, because a freshly produced object is avarpath.
Related
fn add
var fn add(item: Item): Task<Result<Void, Failure>>
Finishes when the target has taken the item.
fn end
var fn end(): Task<Result<Void, Failure>>
No more items. Flushes, reports what is left to report, and makes every further add fail.
fn addAll
var fn addAll(items: Iterate<Item>): Task<Result<Void, Failure>>
Everything an Iterate holds, in order. Does not end the sink: a sink may be filled from several places.
fn fill
var fn fill(var source: Source<Item, Failure>): Task<Result<Void, Failure>>
Reads source empty into this sink and ends it: the same pump as source.into(sink), from this side.
fn mapFailure
var fn mapFailure<Other>(transform: (failure: Failure) => Other): Sink<Item, Other>
A sink for the same items with a failure of another type: what mapError is on a Result.
fn buffered
var fn buffered(capacity: Int = 64): Buffered<Item, Failure>
Collects items and hands them on in groups of capacity, which is what makes many small writes one big one.
The answer is a Buffered, not a Sink, so that flush() stays reachable; finish() flushes.
fn pushing
static fn pushing(accept: (item: Item) => Task<Result<Void, Failure>>, complete: () => Task<Result<Void, Failure>>): Sink<Item, Failure>
Two closures instead of a type: the escape hatch for a sink with no state worth a type of its own. With state,
write an ordinary shared type with var fields.
fn discarding
static fn discarding(): Sink<Item, Failure>
Everything written into it is dropped. What /dev/null is, and what a test writes into.
type Pushing
shared type Pushing<Item, Failure> with Sink<Item, Failure>
Sink.pushing(accept, complete): one closure per verb, and the thing below it that has to be released.
fn of
static fn of(accept: (item: Item) => Task<Result<Void, Failure>>, complete: () => Task<Result<Void, Failure>>, below: Close? = None): Pushing<Item, Failure>
Wraps accept and complete as a sink, and below as what releasing it releases.
fn add
var fn add(item: Item): Task<Result<Void, Failure>>
Calls accept with the item.
fn end
var fn end(): Task<Result<Void, Failure>>
Calls complete.
fn close
var fn close()
Nothing of its own: below is a field, and the release of the sink releases it right after this.
type Buffered
shared type Buffered<Item, Failure> with Sink<Item, Failure>
A sink in front of another one that hands items on in groups. Buffering is a wrapper and never a hidden property of a sink, because "when was it actually written" is the one question a writer has to be able to answer:
var out = standardOutput().buffered(capacity: 4096)
out.addLine(report).await()?
out.flush().await()? // now it is on the screen
out.end().await()? // ... and this flushes too
close() does not flush: it releases, and what was buffered is lost. That is the difference between abandoning
a sink and ending it.
Pitfalls
- If a write below fails partway through a flush, the items after it are dropped, not kept for a retry: a failed sink is over, so there is nothing left to retry into.
Related
Sink.buffered- what answers aBuffered.
fn of
static fn of(var downstream: Sink<Item, Failure>, capacity: Int = 64): Buffered<Item, Failure>
Wraps downstream and buffers up to capacity items before writing a group.
fn add
var fn add(item: Item): Task<Result<Void, Failure>>
Adds to the buffer, and flushes once it reaches capacity.
fn flush
var fn flush(): Task<Result<Void, Failure>>
Writes everything that is buffered, and finishes when it has arrived below.
fn end
var fn end(): Task<Result<Void, Failure>>
Flushes, then ends downstream.
fn close
var fn close()
Nothing of its own, and no flush: whatever is still buffered is lost, and downstream is a field the release of
the sink releases right after this.