Skip to content

Concurrency ​

Talor has no async and no await. Concurrency is two libraries: std.fiber runs many lightweight tasks on a few threads, and std.thread starts operating-system threads. A function that reads a socket has one signature: inside a fiber the read parks the fiber and the thread runs another one, and outside a fiber it blocks the thread. No function is coloured by the way it is called.

What the language contributes is the ownership rules that make it safe: a task receives its values by move, what two tasks share is read-only, and a write two tasks need goes through a lock that owns the data.

Fibers ​

A fiber is a stack of its own and a saved context. run(body) starts the scheduler on the calling thread and returns when every fiber it started has finished. Inside it, spawn(body) starts a fiber and join(f) waits for one.

talor
use std.fiber.{run, spawn, join, sleep_ms};

fn main(): i32 {
    let done = run(() => {
        let slow = spawn(() => {
            match sleep_ms(20) { Ok(_) => {}, Err(_) => {} }
            println("the slow fiber ends second");
        });
        let fast = spawn(() => { println("the fast fiber ends first"); });
        match slow { Ok(f) => { match join(f) { Ok(_) => {}, Err(_) => {} } }, Err(_) => {} }
        match fast { Ok(f) => { match join(f) { Ok(_) => {}, Err(_) => {} } }, Err(_) => {} }
    });
    match done { Ok(_) => 0, Err(_) => 1 }
}

The rest of std.fiber:

CallWhat it does
run_on(n, body)runs the fibers on n threads, the calling one and n - 1 the runtime starts; a ready fiber may be taken by a thread with nothing to run
sleep_ms(ms), yield_now()park the fiber for a time, or let the others run
timeout(ms, body)runs body and answers Err (code timed_out_code()) if it has not finished in ms
cancel(f), cancelled()ask a fiber to stop: it hears the request where it would wait (sleep_ms, yield_now, join), which answers Err there, or by reading cancelled(); nothing stops it from outside, and fibers it started run on
offload(body)runs body on a pool thread and parks the fiber until it returns: the way to call C code that blocks without stopping the other fibers
all(bodies)runs a list of closures at once and answers their values in list order (below)

A fiber's stack does not move, so a view it holds stays valid across a park. A fiber that runs past its stack ends the program with a message naming the fiber and its stack size, instead of writing into another fiber's memory.

all: fibers that borrow what they capture ​

A closure written directly in the list given to fiber.all borrows what it captures for the length of the call, instead of moving it, and the caller has the values back, written, when all returns. Two closures of the list may not touch the same place when one of them writes it; writing two different fields is accepted:

talor
use std.fiber;

struct Halves { left: Array<i64>, right: Array<i64> }

fn main(): i32 {
    let done = fiber.run(() => {
        let t = Halves { left: [1, 2, 3], right: [4, 5, 6] };
        let lens = fiber.all([
            () => { t.left.push(0); t.left.len() },
            () => { t.right.push(0); t.right.len() },
        ]);
        match lens {
            Ok(ls) => println(`lengths ${ls[0]} and ${ls[1]}`),
            Err(e) => println(`failed: ${e.message()}`),
        }
        println(`left ends in ${t.left[3]}, right ends in ${t.right[3]}`);
    });
    match done { Ok(_) => 0, Err(_) => 1 }
}

t.left in one closure and t.right in the other are disjoint; t in one and t.left in the other would be refused, and so would xs[0] and xs[1], because an index is not told apart. A closure in the list may not hand a capture over (return it, store it, or pass it to a parameter that takes it), since the capture is only borrowed. When one returns an error or panics, all answers the first failure in list order, after every closure has finished.

Threads ​

std.thread starts a closure on a thread of its own with start, and join(t) waits for it and detach(t) lets it go. hardware_count() says how many threads the machine runs at once.

The closure given to spawn or start takes its captures by move, so the task owns what it was given. It may be a closure that is called once, because each runs its closure exactly once. A closure that the callee keeps may not capture a view, since the task can outlive the frame the view came from.

Sharing: shared T ​

A value two tasks both need is put in a shared, and each task gets a clone of it: a clone is one more count on the same allocation, not a copy of the value. A place reached through a shared is read-only, because two tasks could write it at the same moment; a write through one is refused at compile time. The count is plain until the value can be reached from a second thread, and atomic from then on, so a program that shares nothing pays nothing for it.

Channels ​

A Channel<T> carries values of any type from one task to another, whole: send takes the value and recv hands it back to the receiver.

talor
use std.thread.{start, join, channel, send, recv, close};

fn main(): i32 {
    let ch = match channel<string>(4) { Ok(c) => c, Err(_) => { return 1; } };
    let owner = shared ch;
    let mine = owner.clone();
    let worker = start(() => {
        for word in ["one", "two", "three"] {
            match send(*mine, word) { Ok(_) => {}, Err(_) => {} }
        }
        match close(*mine) { Ok(_) => {}, Err(_) => {} }
    });
    let going = true;
    while going {
        match recv(*owner) {
            Ok(Some(word)) => println(`received ${word}`),
            Ok(None) => { going = false; }
            Err(_) => { going = false; }
        }
    }
    match worker { Ok(t) => { match join(t) { Ok(_) => 0, Err(_) => 1 } }, Err(_) => 1 }
}
  • channel<T>(capacity) makes one. A positive capacity makes send wait while the channel is full; 0 means no bound.
  • recv waits for a value and answers None once the channel is closed and empty; try_recv never waits.
  • send on a closed channel answers Err, and the value it was given is released there.
  • select(chans, timeout_ms) waits on several channels and answers which one had a value.
  • A channel released with values still in it releases each of them.

Inside a fiber, waiting on a channel parks the fiber; on a thread, it blocks the thread.

Locks own what they guard ​

A lock is not a value standing beside the data: mutex(v) takes v, and the only way to the data is with(m, body), which takes the lock, lends the data to body, and releases the lock when body returns. The body's parameter is a Locked<T>, a place that is writable for the length of that one call and cannot be kept, returned or captured.

talor
use std.thread.{start, join, mutex, with, Thread};

struct Counter { n: i64 }

fn main(): i32 {
    let m = shared mutex(Counter { n: 0 });
    let threads: Array<Thread> = [];
    for i in 0..4 {
        let mine = m.clone();
        match start(() => {
            for k in 0..1000 {
                match with(*mine, (c) => { c.n += 1; }) { Ok(_) => {}, Err(_) => {} }
            }
        }) {
            Ok(t) => threads.push(t),
            Err(_) => { return 1; }
        }
    }
    while threads.len() > 0 {
        match threads.pop() { Some(t) => { match join(t) { Ok(_) => {}, Err(_) => {} } }, None => {} }
    }
    match with(*m, (c) => c.n) {
        Ok(n) => { println(`four threads counted ${n}`); 0 }
        Err(_) => 1,
    }
}
  • try_with(m, body) runs body only if the lock is free, and answers None without waiting otherwise.
  • rwlock(v) admits many readers or one writer: read(l, body) lends the data read-only to any number of readers at once, and write(l, body) lends a Locked<T> to one.
  • A helper can take the lent place by its plain type, fn bump(c: Counter), and be called from the body.
  • A task that takes a lock it already holds is stopped with a panic rather than left waiting forever, and a lock whose holder panicked is poisoned: the next with answers Err.

Atomics ​

For one counter or flag, an Atomic is cheaper than a lock: atomic(n) makes one, and load, store, add and compare_exchange read and write it in one step.

talor
use std.thread.{start, join, atomic, add, load, Thread};

fn main(): i32 {
    let hits = match atomic(0) { Ok(a) => shared a, Err(_) => { return 1; } };
    let threads: Array<Thread> = [];
    for i in 0..4 {
        let mine = hits.clone();
        match start(() => { for k in 0..1000 { add(*mine, 1); } }) {
            Ok(t) => threads.push(t),
            Err(_) => { return 1; }
        }
    }
    while threads.len() > 0 {
        match threads.pop() { Some(t) => { match join(t) { Ok(_) => {}, Err(_) => {} } }, None => {} }
    }
    println(`hits: ${load(*hits)}`);
    0
}

When a task panics ​

A panic ends the fiber or the thread it happens on, and not the process: the task's frames release what they own on the way out, PANIC: and the place are written to standard error, and whoever joins the task gets Err with code 101, the panic's text, and e.panicked() true. A panic on the main thread, or in the body given to run, ends the process with code 101.

talor
use std.fiber.{run, spawn, join};

fn main(): i32 {
    let done = run(() => {
        match spawn(() => { let xs: Array<i64> = []; println(`${xs[3]}`); }) {
            Ok(f) => {
                match join(f) {
                    Ok(_) => println("joined"),
                    Err(e) => println(`the fiber panicked: ${e.panicked()}, code ${e.code()}`),
                }
            }
            Err(_) => {}
        }
        println("the run goes on");
    });
    match done { Ok(_) => 0, Err(_) => 1 }
}

Talor v0.1.0 - Released under the MIT OR Apache-2.0 license.