LEVIATHAN v962456e · 962456eee1

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

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 a RuntimeException.
  • receive() returns a Promise<T?>. Await it to get the next value, or None once the channel is closed and everything sent has been delivered. Test with != None to 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; grow does not exist.
  • A sent value is deep-copied. It must be one of the types allowed across threads (see lang.threads); a Worker, a Promise, a socket or a closure cannot be sent.
  • receive() yields None only after close() and after every sent value has been received.
  • send after close() 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() or std::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() -> void

Close 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) -> void

Add 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.