LEVIATHAN v962456e · 962456eee1

Standard Library

class InStream<T>

The reading end of a stream: a queue of values of type T that something else produces.

since 0.1.0-alpha.1linuxwindows

bases
IDisposable, IIterable<T>

Overview

Streams are how a program hears from the outside world: a Timer hands you an InStream of ticks and signal::on hands you one of signal numbers. The stream holds values that have arrived and not yet been read. You consume them in one of three ways: pull them one at a time with pull or pullOrNone, register a callback with subscribe, or loop over the stream with for. A stream has exactly one consumer, so subscribe and iterator each claim the stream, and a later pull, subscribe or iterator call on a claimed stream throws a RuntimeException.

An InStream is an IDisposable, so using closes it when the scope ends. Closing is safe to repeat and releases whatever the producer attached to the stream, such as a signal subscription. Values pushed into a closed stream are dropped silently. A for loop over a stream finishes when the stream is closed and drained, so a stream that is never closed makes the loop wait for more values without end; use take on asSeq() to bound it.

Description

An InStream<T> is a read view over a StreamBuffer<T>. It offers four ways to take values out, and the first one used becomes the stream's only consumer.

  • pull() returns the next value and throws if there is none.
  • pullOrNone() returns the next value, or None when nothing is queued. It never waits. hasData() tells you whether pull() would succeed.
  • subscribe(callback) hands every value to callback, the ones already queued first and every later one as it is pushed.
  • iterator() is what for (T x in stream) calls. asSeq() joins the lazy Seq<T> pipeline that arrays use, so map, where and take work over a stream too.

Subscribing and the exclusive consumer

StreamBuffer<string> buf = StreamBuffer();
OutStream<string> out = OutStream(buf);
InStream<string> inp = InStream(buf);
out << "queued before";
inp.subscribe((s) => console.writeln("event: ${s}"));
out << "a" << "b";
console.writeln("buffered now: ${buf.count()}");
try {
    inp.pull();
} catch (RuntimeException e) {
    console.writeln(e.message);
}
event: queued before
event: a
event: b
buffered now: 0
consumer end is claimed by a subscriber

Once a callback is subscribed, values go straight to it and nothing stays queued, so a pull() would have nothing to take; it throws instead of returning nothing silently.

Rules

  • subscribe and iterator() are standing claims on the consumer end. The second claim of either kind, and any pull() after a claim, throws a RuntimeException (consumer end is claimed by a subscriber / consumer end is claimed by an iterator). Broadcasting one stream to many readers is something you build by reading one stream and pushing to several.
  • pull() on an empty open stream throws stream is empty. On a closed stream it throws stream is closed, even if values are still queued: close() from the consumer's side means "I am finished; drop everything".
  • pullOrNone(), for loops and asSeq() are the producer-friendly readers: they always hand out whatever is already queued before they report the end of the stream. A producer that writes its last value and then closes the stream (out << v; buf.close();) still has v delivered. pullOrNone() returns None both for "nothing queued yet" and for "closed and empty"; it never waits.
  • close() is idempotent and never throws. It runs any cleanup the producer attached (for a signal::on stream, that is signal::off), then closes the buffer. This is what makes the stream usable with using.
  • A for loop over an open, empty stream suspends the current task until the next value arrives or the stream is closed, then continues. Close the stream from anywhere (the loop's own body, a callback) and the loop ends.
  • A stream that never closes is an infinite source. Terminals such as for or Seq operations that need the whole sequence wait forever on it; bound them with take(n).

Examples

Closing a stream: values written afterwards are dropped, pull() reports the close, and pullOrNone() still drains what was queued.

Closing a stream

StreamBuffer<int> buf = StreamBuffer();
OutStream<int> out = OutStream(buf);
InStream<int> inp = InStream(buf);
out << 7;
inp.close();
out << 8;
console.writeln("count after close: ${buf.count()}");
try {
    inp.pull();
} catch (RuntimeException e) {
    console.writeln(e.message);
}
int? left = inp.pullOrNone();
console.writeln(left ?? -1);
inp.close();
console.writeln("closed twice");

StreamBuffer<int> b2 = StreamBuffer();
OutStream<int> o2 = OutStream(b2);
InStream<int> i2 = InStream(b2);
o2 << 1 << 2;
o2 << 3;
b2.close();
for (int x in i2) { console.writeln("drained ${x}"); }
count after close: 1
stream is closed
7
closed twice
drained 1
drained 2
drained 3

A stream joins the lazy pipeline through asSeq():

A stream as a lazy sequence

IOStream<int> chan = IOStream(StreamBuffer());
chan << 1 << 2 << 3 << 4 << 5;
chan.close();
Seq<int> evens = chan.asSeq().where((n) => n % 2 == 0).map((n) => n * 10);
for (int x in evens) { console.writeln(x); }
20
40

A for loop over a live source delivers every value as it arrives. Here a repeating timer feeds the stream, and the loop ends when the stream is closed from inside the loop:

Iterating a timer's ticks

Timer t = std::every(10);
InStream<int> ticks = t.ticks();
for (int n in ticks) {
    console.writeln("tick ${n}");
    if (n == 3) { t.cancel(); ticks.close(); }
}
console.writeln("loop ended");
tick 1
tick 2
tick 3
loop ended

Notes

InStream<T> has no >> operator: use pull(). Use hasData() before pull() when you cannot tell whether a value is queued, or use pullOrNone() and check for None.

Examples

StreamBuffer<int> buf = StreamBuffer();
OutStream<int> out = OutStream(buf);
InStream<int> inp = InStream(buf);

out << 10 << 20 << 30;
console.writeln(inp.pull());
console.writeln(inp.hasData());
inp.close();
for (int v in inp) {
    console.writeln("drained ${v}");
}
10
true
drained 20
drained 30

Constructors

new

new(StreamBuffer<T> b)

Create the reading end of the stream backed by the buffer b.

Several views over one buffer share its queue: whatever an OutStream over the same buffer writes, this stream reads.

Parameters

b
The buffer the stream reads from.

Examples

StreamBuffer<string> buf = StreamBuffer();
InStream<string> inp = InStream(buf);
OutStream<string> out = OutStream(buf);
out << "hello";
console.writeln(inp.pull());
hello

Methods

asSeq

asSeq() -> Seq<T>

View the stream as a lazy sequence so map, where and take apply to it.

Reading the sequence claims the stream like iterator does. A sequence over a stream that is never closed has no end, so bound it with take before collecting.

Returns

A sequence that reads the stream.

Examples

StreamBuffer<int> buf = StreamBuffer();
OutStream<int> out = OutStream(buf);
InStream<int> inp = InStream(buf);
out << 1 << 2 << 3 << 4;
inp.close();
Array<int> result = inp.asSeq().map((x) => x * 2).where((x) => x > 2).toArray();
console.writeln(result.joinToString(","));
4,6,8

See also: Seq

close

close() -> void

Close the stream, release anything attached to it, and wake a reader that is waiting.

Closing is safe to repeat and never throws. Values written after the close are dropped. Values queued before the close can still be read by pullOrNone or a for loop. When a producer attached cleanup to the stream, such as a signal subscription, it runs once.

Examples

StreamBuffer<int> buf = StreamBuffer();
OutStream<int> out = OutStream(buf);
void readOne() {
    using InStream<int> inp = InStream(buf);
    out << 1;
    console.writeln("inside: ${inp.pull()}");
}
readOne();
out << 2;
InStream<int> again = InStream(buf);
console.writeln("queued after the scope: ${again.hasData()}");
again.close();
again.close();
console.writeln("closed twice without error");
inside: 1
queued after the scope: false
closed twice without error

hasData

hasData() -> bool

Report whether at least one unread value is waiting in the stream.

Returns

true when pull would return a value, false when the queue is empty.

Examples

StreamBuffer<int> buf = StreamBuffer();
OutStream<int> out = OutStream(buf);
InStream<int> inp = InStream(buf);
console.writeln(inp.hasData());
out << 1;
console.writeln(inp.hasData());
false
true

iterator

iterator() -> IIterator<T>

Create an iterator over the values of the stream, which is what a for loop uses.

Asking for the iterator claims the stream, so afterwards pull, subscribe and a second iterator call throw. The iterator hands out queued values first. When the queue is empty and the stream is still open, hasNext waits for the next value to arrive; once the stream is closed and drained it reports false and the loop ends.

Returns

An iterator that reads the stream.

Throws

RuntimeException
when the stream is already claimed.

Examples

StreamBuffer<int> buf = StreamBuffer();
OutStream<int> out = OutStream(buf);
InStream<int> inp = InStream(buf);
out << 7 << 8;
IIterator<int> it = inp.iterator();
try {
    inp.pull();
} catch (RuntimeException e) {
    console.writeln("caught: ${e.message}");
}
inp.close();
while (it.hasNext()) {
    console.writeln(it.next());
}
caught: consumer end is claimed by an iterator
7
8

See also: IIterator

pull

pull() -> T

Take the oldest value out of the stream and return it.

The call never waits. It throws when the stream is empty, when it has been closed (even if values are still queued; use pullOrNone or a for loop to drain those), or when another consumer has claimed it.

Returns

The oldest unread value.

Throws

RuntimeException
when the stream is empty, closed, or claimed by a subscriber or an iterator.

Examples

StreamBuffer<int> buf = StreamBuffer();
OutStream<int> out = OutStream(buf);
InStream<int> inp = InStream(buf);
out << 1 << 2;
console.writeln(inp.pull());
console.writeln(inp.pull());
try {
    inp.pull();
} catch (RuntimeException e) {
    console.writeln("caught: ${e.message}");
}
1
2
caught: stream is empty

See also: pullOrNone

pullOrNone

pullOrNone() -> T | None

Take the oldest value out of the stream, or return None when there is nothing to take.

The call never waits. None means either that nothing has arrived yet or that the stream is closed and fully read. Unlike pull, it still returns values that were queued before the stream was closed.

Returns

The oldest unread value, or None when none is available.

Throws

RuntimeException
when another consumer has claimed the stream.

Examples

StreamBuffer<int> buf = StreamBuffer();
OutStream<int> out = OutStream(buf);
InStream<int> inp = InStream(buf);
out << 7;
console.writeln(inp.pullOrNone());
console.writeln(inp.pullOrNone() == None);
7
true

See also: pull

subscribe

subscribe((T) => void cb) -> void

Register a callback that receives every value of the stream.

Values already queued are delivered to cb immediately, in order, and every later value is delivered as it is written. The call claims the stream: afterwards pull, iterator and a second subscribe throw. Closing the stream stops the deliveries.

Parameters

cb
The function to call with each value.

Throws

RuntimeException
when the stream is closed or already claimed.

Examples

StreamBuffer<int> buf = StreamBuffer();
OutStream<int> out = OutStream(buf);
InStream<int> inp = InStream(buf);
out << 5;
inp.subscribe((v) => console.writeln("saw ${v}"));
out << 6 << 7;
try {
    inp.pull();
} catch (RuntimeException e) {
    console.writeln("caught: ${e.message}");
}
saw 5
saw 6
saw 7
caught: consumer end is claimed by a subscriber

See also

  • OutStream — The writing end of a stream: you push values of type T in, and a reader takes them out.
  • IOStream — A stream that can be both read and written, with both ends over one queue.
  • Seq — A lazy sequence: a pipeline of steps that does no work until its result is requested.
  • Timer — A source of ticks on the event loop.
  • signal — Operating-system signals delivered to the program as streams.
  • IDisposable — The interface of an object that must be cleaned up when its work is finished.