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.bodyis a closure with no parameters, andspawnreturns aWorker<T>whereTis the closure's result type.- A
Worker<T>is aPromise<T>, soawait workeris 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
spawncall. 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; seestd.Channel.std::cpuCount()returns the number of processors that are online, at least1, 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,structvalues,ArrayandMapof such values, and classes whose shape is known statically (with shared parts and cycles preserved). - Some things cannot cross. A nested closure (only the
spawnbody itself may be a closure), an object tied to a descriptor or to the event loop (TcpStream,TcpListener,Timer,Process, and a disposableInStreamsuch as asignal::onsubscription), aBlock, and aWorkerorPromisehandle 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 aChannel<T>, orawaittheWorker<T>thatspawnreturned. - Guard the
spawncall. Abodythat is not a closure throws a catchableRuntimeExceptionat thespawncall, and no worker starts. Atrythat wraps only theawaitdoes not catch it. - Failures arrive at the join. If the body throws, the worker rejects, and
awaitrethrows the failure as aRuntimeExceptioncarrying 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.