Standard Library
class Channel<T>
A queue that carries values from one worker to another.
since 0.1.0-alpha.1linux
Overview
A channel connects a single producer to a single consumer. send copies the value in, receive gives a promise for the next value, and close ends the stream. Create one with a capacity and an overflow policy: Channel(8, std::overflowBlock()). The capacity always comes with a policy for the full case, so there is no unbounded growth by accident.
A channel is the one way to hand values between a worker and the rest of the program. A captured channel handle is shared, not copied: both sides talk to the same queue. On Linux, with workers on separate threads, the queue is a lock-free ring; elsewhere it is an ordinary queue on the one thread. In the ordinary queue the block policy cannot hold a sender back, so it grows past its capacity instead, while drop and error still apply.
Description
A Channel<T> moves values between execution units. You create one with a capacity and an
overflow policy together: Channel(capacity, policy), where policy is one of
std::overflowBlock(): a full channel makes the producer wait (backpressure),std::overflowDrop(): a full channel discards the new value,std::overflowError(): a full channel makessendthrow aRuntimeException.
There is no unbounded policy: a channel always says what happens when it is full.
send(value)copies the value into the channel. Sending on a closed channel throws aRuntimeException.receive()returns aPromise<T?>. Await it to get the next value, orNoneonce the channel is closed and everything sent has been delivered. Test with!= Noneto know whether you got a value.close()says no more values are coming.
A channel handle that a worker captures is a portal, not a copy: the spawning thread and the worker both end up holding the same channel, which is why it is the one thing you can pass to a worker to talk to it. The conduit is for one producer and one consumer.
A worker streaming results through a channel
void run() {
Channel<int> ch = Channel(8, std::overflowBlock());
Worker<int> producer = std::spawn(() => {
for (int i in 1..4) { ch.send(i * i); }
ch.close();
return 3;
});
int count = 0;
int? item = await ch.receive();
while (item != None) {
console.writeln("received ${item}");
count += 1;
item = await ch.receive();
}
console.writeln("channel closed after ${count} items; worker returned ${await producer}");
}
run();
received 1
received 4
received 9
received 16
channel closed after 4 items; worker returned 3
Rules
- Capacity and policy are given together;
growdoes not exist. - A sent value is deep-copied. It must be one of the types allowed across threads (see
lang.threads); aWorker, aPromise, a socket or a closure cannot be sent. receive()yieldsNoneonly afterclose()and after every sent value has been received.sendafterclose()throws.- One producer and one consumer: do not share one channel among several senders or several receivers.
Examples
The three policies at their limits. A drop channel of capacity 2 keeps the first two values; an
error channel throws on the second send; and a closed channel refuses further sends:
Overflow policies
void run() {
Channel<int> dropping = Channel(2, std::overflowDrop());
dropping.send(1);
dropping.send(2);
dropping.send(3);
dropping.close();
int? item = await dropping.receive();
while (item != None) {
console.writeln("kept ${item}");
item = await dropping.receive();
}
Channel<int> strict = Channel(1, std::overflowError());
strict.send(1);
try {
strict.send(2);
} catch (RuntimeException e) {
console.writeln("overflow: ${e.message}");
}
strict.close();
try {
strict.send(3);
} catch (RuntimeException e) {
console.writeln("closed: ${e.message}");
}
}
run();
kept 1
kept 2
overflow: channel overflow
closed: send on a closed channel
Examples
A worker feeding the main program through a channel
Channel<int> ch = Channel(8, std::overflowBlock());
Worker<int> producer = std::spawn(() => {
for (int i in 1..3) ch.send(i * i);
ch.close();
return 3;
});
int sent = await producer;
console.writeln("sent ${sent}");
int? item = await ch.receive();
while (item != None) {
console.writeln("received ${item}");
item = await ch.receive();
}
console.writeln("closed and drained");
sent 3
received 1
received 4
received 9
closed and drained
Constructors
new
new(int capacity, int p)Create a channel with a capacity and an overflow policy.
Parameters
- capacity
- The number of values the channel holds before the policy applies.
- p
- What to do when the channel is full:
std::overflowBlock(),std::overflowDrop()orstd::overflowError().
Examples
Channel<string> ch = Channel(4, std::overflowBlock());
ch.send("hello");
ch.send("world");
console.writeln(await ch.receive() ?? "none");
console.writeln(await ch.receive() ?? "none");
hello
world
Methods
close
close() -> voidClose the channel.
No more values can be sent. Values already queued can still be received; after them, receive gives None.
Examples
Channel<int> ch = Channel(4, std::overflowBlock());
ch.send(1);
ch.close();
console.writeln(await ch.receive() ?? -1);
int? last = await ch.receive();
console.writeln(last == None);
1
true
receive
receive() -> Promise<T | None>Get a promise for the next value.
The promise resolves with the next value when one is available. Once the channel has been closed and every queued value has been received, it resolves with None, so a loop that ends at None reads everything exactly once. Narrow the result with != None, or use ??.
Returns
A promise for the next value, or for None when the channel is closed and empty.
Examples
Channel<int> ch = Channel(4, std::overflowBlock());
ch.send(5);
ch.send(6);
ch.close();
int? v = await ch.receive();
while (v != None) {
console.writeln("got ${v}");
v = await ch.receive();
}
console.writeln("end of channel");
got 5
got 6
end of channel
send
send(T value) -> voidAdd a value to the channel.
The value is copied, so the sender can keep changing the original. If a receiver is already waiting it gets the value directly. If the channel is full the overflow policy decides: the sender waits, the value is dropped, or an exception is thrown. Sending on a closed channel always throws.
Parameters
- value
- The value to send.
Throws
RuntimeException- when the channel is closed, or when it is full and its policy is
std::overflowError().
Examples
Channel<int> ch = Channel(1, std::overflowError());
ch.send(1);
try {
ch.send(2);
} catch (RuntimeException e) {
console.writeln("caught: ${e.message}");
}
Channel<int> closed = Channel(1, std::overflowBlock());
closed.close();
try {
closed.send(3);
} catch (RuntimeException e) {
console.writeln("caught: ${e.message}");
}
caught: channel overflow
caught: send on a closed channel
See also
- Worker — The handle to a value that another worker is computing.
- overflowBlock — The overflow policy that makes a full channel hold the sender back.
- Promise — A value that arrives later.