Axle v0.14.1

Concurrency primitives

These primitives coordinate threads in the thread world (spawn / Task<T>). If you’re in the async world (async fn / Future<T>) and never sharing state between OS threads, you usually don’t need any of this — the event loop is single-threaded, so coroutines never race with each other.

std/concurrent covers : atomic counters, locks of various flavours, channels for message passing, and rendezvous primitives. See Multithreading for the two-worlds model and how to spawn tasks in the first place.

When to use what

You need to …Reach for …
share a single int / long / bool, atomic ops onlyAtomicI32 / AtomicI64 / AtomicBoolean
protect a multi-step read-modify-writesynchronized (lock) { … }
read often, write rarelyReadWriteLock
pass values between threadsBoundedChannel / UnboundedChannel
wait until N tasks finishCountDownLatch
cap concurrent in-flight workSemaphore
explicit reentrant lock objectReentrantLock

Atomic primitives

AtomicI32 — lock-free counter

use std::concurrent::AtomicI32;

fn main() : i32 {
    let counter : AtomicI32 = new AtomicI32(0);
    counter.incrementAndGet();                 // 0 → 1, returns the new value
    counter.addAndGet(10);                     // returns the new value (11)
    counter.decrementAndGet();                 // back to 10, returns it
    counter.set(0);
    if (counter.get() != 0) { return 1; }
    return 0;
}

AtomicI32 exposes :

  • get() / set(v) — atomic load / store
  • incrementAndGet() / decrementAndGet() — add / subtract one, each returning the new value (there is no bare increment() / decrement(); the value is always returned so the caller never needs a separate read)
  • addAndGet(d) — returns the new value after adding d
  • compareAndSet(expected, update) : bool — stores update only if the value currently equals expected; true when the swap happened
  • swap(newValue) — stores newValue and returns the value it replaced

Cross-thread example :

fn bump1000(hits : AtomicI32) : i32 {
    let i : i32 = 0;
    while (i < 1000) { hits.incrementAndGet(); i = i + 1; }
    return 0;
}

fn main() : i32 {
    let hits : AtomicI32 = new AtomicI32(0);

    let f1 : Task<i32> = spawn bump1000(hits);
    let f2 : Task<i32> = spawn bump1000(hits);

    f1.join();
    f2.join();
    return hits.get();                         // 2000 — never less, never more
}

AtomicI64 — same shape, 64-bit

use std::concurrent::AtomicI64;

fn totalAfterAdd() : i64 {
    let bytes : AtomicI64 = new AtomicI64(0);
    bytes.addAndGet(4096);
    return bytes.get();
}

Identical method set to AtomicI32 — get() / set(v) / incrementAndGet() / decrementAndGet() / addAndGet(d) / compareAndSet(expected, update) / swap(newValue) — only the width differs (64-bit values, i64 arguments and results).

AtomicBoolean — flag with CAS

use std::concurrent::AtomicBoolean;

fn runOnceInitialization() : void {
    println("initialised");
}

fn main() : i32 {
    let started : AtomicBoolean = new AtomicBoolean(false);

    // Idempotent init — only the first caller does the work.
    if (started.compareAndSet(false, true)) {
        runOnceInitialization();
    }
    return 0;
}

compareAndSet(expected, update) returns true if the swap fired, false if the current value didn’t match expected. Standard CAS semantics.

synchronized blocks

The simplest mutual-exclusion primitive: pass any object, take its identity-keyed reentrant lock for the block.

class Account {
    pub(file) id      : i32;
    pub(file) balance : i32;
    constructor(id : i32, initial : i32) {
        self.id = id;
        self.balance = initial;
    }
}

fn transfer(from : Shared<Account>, to : Shared<Account>, amount : i32) : void {
    // Lock ordering by a stable key (the account id) avoids AB-BA
    // deadlocks : every thread takes the two locks in the same order.
    if (from.id < to.id) {
        synchronized (from) {
            synchronized (to) {
                from.balance = from.balance - amount;
                to.balance   = to.balance + amount;
            }
        }
    } else {
        synchronized (to) {
            synchronized (from) {
                from.balance = from.balance - amount;
                to.balance   = to.balance + amount;
            }
        }
    }
}

Properties :

  • Re-entrant : a thread holding the lock can re-enter without deadlocking.
  • Exception-safe : the lock releases even if the body throws.
  • Block-scoped : no try { lock.lock(); … } finally { lock.unlock(); } boilerplate.

The lock-less form shares one global lock

synchronized { … } without a parenthesised object also compiles. It has no key, so every lock-less synchronized block in the whole program takes the same lock — two unrelated critical sections serialise against each other:

use std::concurrent::AtomicI32;

fn main() : i32 {
    let counter : AtomicI32 = new AtomicI32(0);
    synchronized {          // no key: one process-wide lock
        counter.incrementAndGet();
    }
    return counter.get();
}

Use it only for a genuine “one at a time, program-wide” section. Whenever two critical sections protect different state, name the object (synchronized (from) { … }) so each gets its own lock.

ReentrantLock — explicit lock object

When you need tryLock semantics or want a lock object you can pass around explicitly :

use std::concurrent::ReentrantLock;

fn doProtectedWork() : void {
    println("in the critical section");
}

fn main() : i32 {
    let lock : ReentrantLock = new ReentrantLock();

    if (lock.tryLock()) {
        defer lock.unlock();
        doProtectedWork();
    } else {
        // lock was busy — back off
    }
    return 0;
}

Methods :

  • lock() — block until acquired
  • tryLock() — non-blocking, returns bool
  • unlock() — release ; only the holder may call, and a call from any other thread is refused rather than obeyed
  • isLocked() — query, no synchronization (advisory)

Pair lock() / tryLock() with defer unlock() immediately.

ReadWriteLock — multi-reader single-writer

When reads dominate writes :

use std::concurrent::ReadWriteLock;
use std::collections::HashMap;

class Cache {
    rw    : ReadWriteLock;
    data  : HashMap<i32, i32>,

    constructor() {
        self.rw   = new ReadWriteLock();
        self.data = new HashMap<i32, i32>();
    }

    pub fn lookup(self, key : i32) : i32 {
        self.rw.readLock();
        defer self.rw.readUnlock();
        return self.data.getOrDefault(key, -1);   // -1 on a miss
    }

    pub fn store(self, key : i32, value : i32) : void {
        self.rw.writeLock();
        defer self.rw.writeUnlock();
        self.data.put(key, value);
    }
}

Why i32 keys, not string? Inserting into a HashMap moves the key and value into the container. A string parameter is only borrowed from the caller, so forwarding it straight into put would try to give away something you don’t own (rejected as E0504). Scalar keys like i32 are copied, so they forward freely — see Shared<T> and ownership for the moving rules.

Many readers can hold readLock() simultaneously ; writeLock() excludes all readers and other writers. Not re-entrant across the read/write boundary — a thread holding readLock must release it before calling writeLock.

Channels

Channels carry i32 values between threads. The bounded variant applies backpressure ; the unbounded variant grows the buffer indefinitely.

BoundedChannel

spawn takes a function call, so the producer and consumer are named functions :

use std::concurrent::BoundedChannel;

fn produce(ch : BoundedChannel) : i32 ! InterruptedException {
    ch.send(1);
    ch.send(2);
    ch.send(3);                                // blocks until consumer drains
    return 0;
}

fn consume(ch : BoundedChannel) : i32 ! InterruptedException {
    let sum : i32 = 0;
    for (i of 0..3) { sum = sum + ch.receive(); }
    return sum;
}

fn main() : i32 ! InterruptedException, IOException {
    let ch : BoundedChannel = new BoundedChannel(2);

    let producer : Task<i32> = spawn produce(ch);
    let consumer : Task<i32> = spawn consume(ch);

    producer.join();
    let result : i32 = consumer.join();          // 6

    ch.close();
    return result;
}

Non-blocking variants

use std::concurrent::BoundedChannel;

fn main() : i32 ! IOException {
    let ch : BoundedChannel = new BoundedChannel(1);

    if (ch.trySend(42)) {
        // value queued
    } else {
        // channel full — drop, retry, or back off
    }

    let v : i32 = ch.tryReceive();             // 0 if empty (sentinel)
    ch.close();
    return v - 42;
}

trySend returns false on a full channel ; tryReceive returns 0 on an empty one. The sentinel is indistinguishable from a queued zero, so reserve tryReceive for protocols where 0 is never a legitimate payload.

UnboundedChannel

Same API minus the capacity argument :

use std::concurrent::UnboundedChannel;

fn main() : i32 ! InterruptedException, IOException {
    let ch : UnboundedChannel = new UnboundedChannel();
    ch.send(1); ch.send(2); ch.send(3);        // never blocks
    ch.close();
    return 0;
}

Buffer grows without bound — only use when you can prove the producer is rate-limited elsewhere.

CountDownLatch — single-shot barrier

use std::concurrent::CountDownLatch;

fn initializeWorker() : void {
    println("worker ready");
}

fn startAcceptingRequests() : void {
    println("accepting");
}

fn initWorker(ready : CountDownLatch) : i32 {
    initializeWorker();
    ready.countDown();
    return 0;
}

fn main() : i32 ! InterruptedException {
    let ready : CountDownLatch = new CountDownLatch(3);

    let i : i32 = 0;
    while (i < 3) {
        spawn initWorker(ready);
        i = i + 1;
    }

    ready.waitFor();                            // unblocks once all 3 fire
    startAcceptingRequests();
    return 0;
}

waitFor() blocks until the counter hits zero ; waitForWithTimeout(ms) returns false on timeout instead of blocking forever. The counter is single-use — re-arm with a fresh latch.

Semaphore — permit pool

Cap the number of concurrent operations :

use std::concurrent::Semaphore;

fn handleRequest() : void {
    println("served");
}

fn handleOne(connPool : Semaphore) : i32 ! InterruptedException {
    connPool.acquire();
    defer connPool.release();
    handleRequest();
    return 0;
}

fn main() : i32 ! InterruptedException {
    let connPool : Semaphore = new Semaphore(8);

    let i : i32 = 0;
    while (i < 100) {
        spawn handleOne(connPool);
        i = i + 1;
    }
    return 0;
}

At most 8 worker threads run handleRequest simultaneously ; the rest park inside acquire(). tryAcquire() is the non-blocking variant.

Worker pools — spawn + a channel work queue

std::concurrent declares a ThreadPoolExecutor class (and a Thread class), but the callback path that would let the pool invoke a Runnable is not implemented — execute(…) raises UnsupportedOperationException, as does Thread.start(). Build a long-lived worker pool out of the primitives that do work : spawn for the workers, a channel as the work queue, and a sentinel value for shutdown :

use std::concurrent::BoundedChannel;

fn workerLoop(jobs : BoundedChannel) : i32 ! InterruptedException {
    let done : i32 = 0;
    let job : i32 = jobs.receive();
    while (job != -1) {                  // -1 is the shutdown sentinel
        // … handle the job …
        done = done + 1;
        job = jobs.receive();
    }
    return done;
}

fn main() : i32 ! InterruptedException, IOException {
    let jobs : BoundedChannel = new BoundedChannel(8);

    let pool : Task<i32>[4];
    for (w of 0..4) {
        pool[w] = spawn workerLoop(jobs);
    }

    for (job of 0..20) {                 // enqueue the work
        jobs.send(job);
    }
    for (s of 0..4) {                    // one sentinel per worker
        jobs.send(-1);
    }

    let handled : i32 = 0;
    for (w of 0..4) {
        handled = handled + (pool[w].join());
    }
    jobs.close();                        // Closeable — close once drained
    return handled;                      // 20
}

Each worker loops on receive() until it reads the sentinel ; main enqueues one sentinel per worker, then awaits each worker’s processed-job count. See the std/concurrent reference for the ThreadPoolExecutor surface and its current limitation.

Choosing between synchronized and explicit locks

synchronized (X) { … }      // block-scoped, exception-safe, reentrant
ReentrantLock                // explicit, can pass around, tryLock available
ReadWriteLock                // multi-reader / single-writer

Default to synchronized — less ceremony, automatic release, and the lock travels with the data. Reach for ReentrantLock when you need tryLock, and ReadWriteLock when reads vastly outnumber writes.

See also

concurrencythreadsmutexchannelsatomics