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 only | AtomicI32 / AtomicI64 / AtomicBoolean |
| protect a multi-step read-modify-write | synchronized (lock) { … } |
| read often, write rarely | ReadWriteLock |
| pass values between threads | BoundedChannel / UnboundedChannel |
| wait until N tasks finish | CountDownLatch |
| cap concurrent in-flight work | Semaphore |
| explicit reentrant lock object | ReentrantLock |
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 / storeincrementAndGet()/decrementAndGet()— add / subtract one, each returning the new value (there is no bareincrement()/decrement(); the value is always returned so the caller never needs a separate read)addAndGet(d)— returns the new value after addingdcompareAndSet(expected, update) : bool— storesupdateonly if the value currently equalsexpected;truewhen the swap happenedswap(newValue)— storesnewValueand 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 acquiredtryLock()— non-blocking, returnsboolunlock()— release ; only the holder may call, and a call from any other thread is refused rather than obeyedisLocked()— 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
i32keys, notstring? Inserting into aHashMapmoves the key and value into the container. Astringparameter is only borrowed from the caller, so forwarding it straight intoputwould try to give away something you don’t own (rejected as E0504). Scalar keys likei32are copied, so they forward freely — seeShared<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
- Multithreading basics — the two worlds (
spawn/join,async/await) Shared<T>reference counting — lifetimestd/concurrentreference — full API- Concept index — every concurrency primitive cross-linked