LEVIATHAN v962456e · 962456eee1

Library

Threads, workers and message passing

Run work in parallel with std::spawn; values cross between threads by copy, and Channel is the way to talk to a running worker.

since 0.1.0-alpha.1linux

Description

A worker is a piece of work that runs on another thread of execution. The model is isolation with copying: a worker captures its inputs by copy, runs, and returns a result, and every value that crosses a thread boundary is deep-copied. No two threads ever share the same object, so there are no data races and no locks to write.

  • std::spawn(body) starts a worker. body is a closure with no parameters, and spawn returns a Worker<T> where T is the closure's result type.
  • A Worker<T> is a Promise<T>, so await worker is the join: it waits for the result and rethrows any failure. There is no second way to wait.
  • The variables the closure captures are snapshotted at the spawn call. Changing the original afterwards is not seen by the worker, and changes the worker makes to its copy do not reach the original.
  • Channel<T> carries values to or from a worker while it runs; see std.Channel.
  • std::cpuCount() returns the number of processors that are online, at least 1, for sizing a pool of workers.

Parallel sums and the spawn-time snapshot

int sumRange(int from, int to) {
    int total = 0;
    for (int i in from..to) { total += i; }
    return total;
}

Worker<int> low = std::spawn(() => sumRange(1, 500));
Worker<int> high = std::spawn(() => sumRange(501, 1000));
int a = await low;
int b = await high;
console.writeln("low=${a} high=${b} total=${a + b}");
console.writeln("cpuCount is at least 1: ${std::cpuCount() >= 1}");
low=125250 high=375250 total=500500
cpuCount is at least 1: true

Rules

  • Copy always. What crosses a thread boundary is copied by flattening and rebuilding. The values that can cross are: numbers, char, bool, None, strings, ranges, struct values, Array and Map of such values, and classes whose shape is known statically (with shared parts and cycles preserved).
  • Some things cannot cross. A nested closure (only the spawn body itself may be a closure), an object tied to a descriptor or to the event loop (TcpStream, TcpListener, Timer, Process, and a disposable InStream such as a signal::on subscription), a Block, and a Worker or Promise handle are all rejected with a catchable error that names the type. Each worker opens its own sockets and timers. To talk to a worker, pass a Channel<T>, or await the Worker<T> that spawn returned.
  • Guard the spawn call. A body that is not a closure throws a catchable RuntimeException at the spawn call, and no worker starts. A try that wraps only the await does not catch it.
  • Failures arrive at the join. If the body throws, the worker rejects, and await rethrows the failure as a RuntimeException carrying the original message.
  • Let the spawning thread print. The deterministic way to write worker programs is: workers compute and return or send; only the spawning thread prints, after the join. Printing from a worker body races with other output.
  • Pass inputs as locals. The snapshot covers captured locals and parameters. A top-level variable of the program is shared state rather than a capture, so a worker reads whatever it holds when it runs; put a worker's inputs in locals instead.
  • Real parallelism needs the native backend. On the compiled backend a worker is a real operating-system thread with its own heap and event loop; the results are identical to the cooperative way the interpreters run workers one after another.

Examples

The snapshot is taken when spawn is called. Here the worker computes with 10 and a box whose value was 10, even though the originals change right after:

Captures are copied at spawn

class Box { int value = 10; }

void demo() {
    Box box = Box();
    int limit = 10;
    Worker<int> w = std::spawn(() => limit * box.value);
    limit = 99;
    box.value = 5;
    console.writeln("worker computed ${await w}");
    console.writeln("original box still changed: ${box.value}");
}
demo();
worker computed 100
original box still changed: 5

Notes

Immutable data could be shared between threads without copying, but sharing is not part of the model.

See also

  • Worker — The handle to a value that another worker is computing.
  • Channel — A queue that carries values from one worker to another.
  • Promise — A value that arrives later.
  • TaskGroup — A set of tasks that live and end together.