Reference

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 await on add. 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. add must 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. After end(), add fails. The word is not finish, which belongs to Accumulator and answers a result: the end of a stream is not a value.
  • A failure ends the stream. After a failure neither add nor end will succeed again.
  • close() is the abrupt end: it releases the target and cannot fail, which is what Close needs: the last release of the sink runs it, and no program calls it. A sink that is released without end() may have written less than it was given - that is the point of the two.
  • Buffering is never implicit. sink.buffered(capacity:) says it, and Buffered.flush() is where it is forced.
  • Whoever writes needs the permission - now or later. mapFailure and buffered take var self as well, because the wrapper keeps the sink in a var field and writes into it afterwards. A chain is still one expression, because a freshly produced object is a var path.

Related

  • Source - the reading end, and into/fill's other side.
  • Buffered - what buffered(capacity:) answers.

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

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.