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, orNonewhen nothing is queued. It never waits.hasData()tells you whetherpull()would succeed.subscribe(callback)hands every value tocallback, the ones already queued first and every later one as it is pushed.iterator()is whatfor (T x in stream)calls.asSeq()joins the lazySeq<T>pipeline that arrays use, somap,whereandtakework 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
subscribeanditerator()are standing claims on the consumer end. The second claim of either kind, and anypull()after a claim, throws aRuntimeException(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 throwsstream is empty. On a closed stream it throwsstream is closed, even if values are still queued:close()from the consumer's side means "I am finished; drop everything".pullOrNone(),forloops andasSeq()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 hasvdelivered.pullOrNone()returnsNoneboth 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 asignal::onstream, that issignal::off), then closes the buffer. This is what makes the stream usable withusing.- A
forloop 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
fororSeqoperations that need the whole sequence wait forever on it; bound them withtake(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() -> voidClose 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() -> boolReport 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() -> TTake 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 | NoneTake 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) -> voidRegister 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
Tin, 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.