Learn Rust Series (#72) - Condvars, Barriers, and Thread Coordination
Learn Rust Series (#72) - Condvars, Barriers, and Thread Coordination
What will I learn
- You will learn how a condition variable lets a thread sleep until another signals a change;
- why a Condvar always pairs with a
Mutex, and why you wait inside a loop (spurious wakeups); - how
wait_whileexpresses "sleep until this predicate is false" cleanly; - how a Barrier makes a fixed group of threads rendezvous before any proceeds;
- how to build a producer/consumer and a count-down latch from these primitives.
Requirements
- A working modern computer running macOS, Windows or Ubuntu;
- An installed Rust toolchain (via rustup, from rustup.rs);
- The previous seventy-one episodes, especially
Mutex,Arc, and threads; - The ambition to learn systems programming from the ground up.
Difficulty
- Advanced
Curriculum (of the Learn Rust Series):
- Learn Rust Series (#1) - Introduction to Rust
- Learn Rust Series (#2) - Variables, Types, Functions
- Learn Rust Series (#3) - Ownership & Borrowing
- Learn Rust Series (#4) - Control Flow & Pattern Matching
- Learn Rust Series (#5) - Structs & Enums
- Learn Rust Series (#6) - Error Handling
- Learn Rust Series (#7) - Collections
- Learn Rust Series (#8) - Traits & Generics
- Learn Rust Series (#9) - Modules & Crates
- Learn Rust Series (#10) - Lifetimes
- Learn Rust Series (#11) - Closures & the Iterator Trait
- Learn Rust Series (#12) - Smart Pointers: Box, Rc & RefCell
- Learn Rust Series (#13) - Concurrency: Threads, Channels, Arc & Mutex
- Learn Rust Series (#14) - Mini Project: A Command-Line To-Do App
- Learn Rust Series (#15) - Trait Objects & Dynamic Dispatch
- Learn Rust Series (#16) - Static vs Dynamic Dispatch
- Learn Rust Series (#17) - Associated Types vs Generic Parameters
- Learn Rust Series (#18) - Operator Overloading with std::ops
- Learn Rust Series (#19) - Deref, DerefMut & Deref Coercion
- Learn Rust Series (#20) - Drop & Deterministic Destruction (RAII)
- Learn Rust Series (#21) - From, Into, TryFrom & Idiomatic Conversions
- Learn Rust Series (#22) - Deriving Common Traits
- Learn Rust Series (#23) - The Orphan Rule & Trait Coherence
- Learn Rust Series (#24) - Blanket Implementations & the Newtype Pattern
- Learn Rust Series (#25) - Marker Traits: Sized, Send, Sync & Copy
- Learn Rust Series (#26) - Const Generics: Types That Depend on Values
- Learn Rust Series (#27) - Generic Associated Types & Lending Iterators
- Learn Rust Series (#28) - Sealed Traits & Designing Stable APIs
- Learn Rust Series (#29) - Typestate Programming: State Machines in the Type System
- Learn Rust Series (#30) - Mini Project: A Generic Units-of-Measure Library
- Learn Rust Series (#31) - Move Semantics Deep Dive
- Learn Rust Series (#32) - Interior Mutability: Cell & RefCell
- Learn Rust Series (#33) - Rc Internals: Reference Counting & Shared Ownership
- Learn Rust Series (#34) - Arc: Thread-Safe Reference Counting & Its Cost
- Learn Rust Series (#35) - Weak References & Breaking Reference Cycles
- Learn Rust Series (#36) - Cow: Clone-on-Write for Borrow-or-Own APIs
- Learn Rust Series (#37) - Pin & Self-Referential Structs
- Learn Rust Series (#38) - PhantomData, Zero-Sized Types & Marker Lifetimes
- Learn Rust Series (#39) - Variance: Covariance, Contravariance & Why It Matters
- Learn Rust Series (#40) - Arena & Bump Allocation Patterns
- Learn Rust Series (#41) - Building Your Own Smart Pointer
- Learn Rust Series (#42) - Drop Order, the Drop Check & Leak Safety
- Learn Rust Series (#43) - std::mem: swap, replace, take & forget
- Learn Rust Series (#44) - Higher-Ranked Trait Bounds & Lifetime Elision
- Learn Rust Series (#45) - Mini Project: A Doubly-Linked List, Safe then Unsafe
- Learn Rust Series (#46) - Result Combinators: map, map_err, and_then, ok_or
- Learn Rust Series (#47) - Option Combinators & Null-Free Programming
- Learn Rust Series (#48) - Custom Error Types & the std::error::Error Trait
- Learn Rust Series (#49) - thiserror: Ergonomic Library Errors
- Learn Rust Series (#50) - anyhow: Flexible Application-Level Errors & Context
- Learn Rust Series (#51) - Panics, Unwinding, abort, and catch_unwind
- Learn Rust Series (#52) - Testing: Unit Tests, Integration Tests, and Doctests
- Learn Rust Series (#53) - Property-Based Testing with proptest
- Learn Rust Series (#54) - Fuzzing with cargo-fuzz and libFuzzer
- Learn Rust Series (#55) - Benchmarking with Criterion and Reading the Numbers
- Learn Rust Series (#56) - Cargo Workspaces and Multi-Crate Projects
- Learn Rust Series (#57) - Feature Flags and Conditional Compilation (cfg)
- Learn Rust Series (#58) - Build Scripts (build.rs) and Generating Code at Build Time
- Learn Rust Series (#59) - Clippy, rustfmt, and Writing Idiomatic Rust
- Learn Rust Series (#60) - Mini Project: A Fully Tested, Documented, Published-Ready CSV Toolkit Crate
- Learn Rust Series (#61) - Send and Sync: The Traits Behind Fearless Concurrency
- Learn Rust Series (#62) - Scoped Threads: Borrowing Local Data Across Threads
- Learn Rust Series (#63) - Channels: mpsc, Ownership Transfer, and Backpressure
- Learn Rust Series (#64) - Crossbeam: Faster Channels and Scoped Concurrency
- Learn Rust Series (#65) - Mutex, RwLock, and Handling Lock Poisoning
- Learn Rust Series (#66) - Atomics: AtomicUsize, fetch_add, and Compare-and-Swap
- Learn Rust Series (#67) - Memory Ordering: Relaxed, Acquire, Release, and SeqCst
- Learn Rust Series (#68) - Building a Lock-Free Stack (Treiber)
- Learn Rust Series (#69) - Building a Lock-Free Queue (Michael-Scott)
- Learn Rust Series (#70) - Data Parallelism with Rayon and Parallel Iterators
- Learn Rust Series (#71) - Custom Thread Pools and Work Stealing
- Learn Rust Series (#72) - Condvars, Barriers, and Thread Coordination (this post)
Learn Rust Series (#72) - Condvars, Barriers, and Thread Coordination
Last episode we built a thread pool, and its whole design rested on a channel: workers block inside recv() until a job arrives, and the channel wakes them at exactly the right moment. But look one level down and ask -- how does recv() know when to wake up? It does not poll a flag in a tight loop, burning a CPU core while it waits. Somewhere underneath there is a primitive that lets a thread genuinely go to sleep, hand its core back to the operating system, and be woken the instant another thread changes the state it cares about. That primitive is the condition variable, or Condvar, and it is the topic of today.
A Mutex protects data, but a Mutex alone cannot express waiting for a condition that another thread will bring about: "sleep until the queue is non-empty", "sleep until all workers have finished phase one", "sleep until this flag flips to true". You could, of course, lock the mutex, check, unlock, and spin around again -- a busy-loop -- but that pins a core at 100% doing nothing useful. The Condvar is the efficient answer: it lets a thread release its lock, sleep, and be woken precisely when something changes. Paired with a Barrier for group rendezvous, these are the coordination primitives that turn a bag of independent threads into a choreographed system ;-)
Solutions to Episode 71 Exercises
Episode 71 was custom thread pools and work stealing. Here are worked solutions to the three exercises.
Exercise 1 -- a ThreadPool with a new(size) constructor and an execute method, running ten jobs and shutting down cleanly when dropped:
use std::sync::{mpsc, Arc, Mutex};
use std::thread;
type Job = Box<dyn FnOnce() + Send + 'static>;
struct ThreadPool {
workers: Vec<thread::JoinHandle>,
sender: Option<mpsc::Sender>,
}
impl ThreadPool {
fn new(size: usize) -> Self {
let (sender, receiver) = mpsc::channel::();
let receiver = Arc::new(Mutex::new(receiver));
let mut workers = Vec::with_capacity(size);
for _ in 0..size {
let receiver = Arc::clone(&receiver);
workers.push(thread::spawn(move || {
while let Ok(job) = receiver.lock().unwrap().recv() {
job();
}
}));
}
ThreadPool { workers, sender: Some(sender) }
}
fn execute<F: FnOnce() + Send + 'static>(&self, f: F) {
self.sender.as_ref().unwrap().send(Box::new(f)).unwrap();
}
}
impl Drop for ThreadPool {
fn drop(&mut self) {
drop(self.sender.take()); // close the channel so workers leave their loop
for w in self.workers.drain(..) {
w.join().unwrap();
}
}
}
fn main() {
let pool = ThreadPool::new(4);
for i in 0..10 {
pool.execute(move || println!("job {i} ran"));
}
// pool dropped here: pending jobs drain, channel closes, every worker joins
}
The key insight is the Option<Sender> plus the Drop: taking and dropping the sender closes the channel, so each worker's recv() returns Err, every loop ends, and the joins return instead of hanging forever.
Exercise 2 -- each job sends its result into a channel, and after submitting we drop our own sender and collect a sorted Vec<i32>:
use std::sync::mpsc;
use std::thread;
fn main() {
let (jtx, jrx) = mpsc::channel::<i32>();
let (rtx, rrx) = mpsc::channel::<i32>();
let rx = std::sync::Arc::new(std::sync::Mutex::new(jrx));
let mut workers = Vec::new();
for _ in 0..3 {
let (rx, rtx) = (std::sync::Arc::clone(&rx), rtx.clone());
workers.push(thread::spawn(move || {
while let Ok(n) = rx.lock().unwrap().recv() {
rtx.send(n * n).unwrap();
}
}));
}
drop(rtx); // drop our copy so the result channel can close
for i in 1..=5 { jtx.send(i).unwrap(); }
drop(jtx);
let mut out: Vec<i32> = rrx.iter().collect();
for w in workers { w.join().unwrap(); }
out.sort();
println!("{out:?}"); // [1, 4, 9, 16, 25]
}
Two closes carry the correctness: drop(jtx) ends the workers' input loops, and drop(rtx) (our own copy of the result sender) lets rrx.iter() terminate once the last worker's clone is gone. Forget either drop and the program hangs.
Exercise 3 -- why holding the receiver's MutexGuard across job() would serialise the pool:
// A worker's loop is: while let Ok(job) = receiver.lock().unwrap().recv() { job(); }
//
// The temporary MutexGuard from receiver.lock() is created to call recv(), and
// because it is never bound to a variable, it is DROPPED at the end of the
// `while let` condition -- BEFORE job() runs. So a worker holds the lock only
// long enough to grab a job, then releases it, freeing every other worker to
// grab their own job concurrently.
//
// If instead you wrote:
// let guard = receiver.lock().unwrap();
// while let Ok(job) = guard.recv() { job(); } // guard held across job()!
// the guard would live for the WHOLE loop. One worker would own the lock for the
// entire run, every other worker would block on receiver.lock(), and jobs would
// execute one at a time. The "pool" would be a single thread in disguise.
The whole point of a pool is parallel execution, so the lock must cover only the hand-off, never the work.
The condition variable
A Condvar never travels alone -- it always pairs with a Mutex, conventionally as a (Mutex<T>, Condvar) tuple wrapped in an Arc for sharing. The Mutex guards the state; the Condvar handles the waiting on that state. A waiting thread locks, checks the condition, and calls wait, which does something that would be impossible to write by hand safely: it atomically releases the lock and puts the thread to sleep. A signalling thread later locks, updates the data, and calls notify_one to wake a sleeper:
use std::sync::{Arc, Mutex, Condvar};
use std::thread;
fn main() {
let pair = Arc::new((Mutex::new(false), Condvar::new()));
let p2 = Arc::clone(&pair);
thread::spawn(move || {
let (lock, cvar) = &*p2;
*lock.lock().unwrap() = true; // change the state
cvar.notify_one(); // then wake the waiter
});
let (lock, cvar) = &*pair;
let mut ready = lock.lock().unwrap();
while !*ready {
ready = cvar.wait(ready).unwrap(); // release the lock and sleep until notified
}
println!("signalled: {}", *ready); // signalled: true
}
The atomicity is the entire reason a Condvar exists. Imagine you tried to do this yourself: unlock the mutex, then go to sleep. Between those two steps the signalling thread could sneak in, set the flag, and fire its notification -- into the void, because you are not asleep yet. You would then sleep forever, having missed the one wake-up you needed. This is the lost-wakeup race, and wait closes it by making "release the lock" and "start sleeping" a single indivisible operation. When wait returns, it has also re-acquired the lock for you, so *ready is safe to read again.
Always wait in a loop
Notice the while !*ready around the wait -- not an if. This is mandatory, and it trips up almost everyone the first time. A Condvar is permitted to wake spuriously: wait can return even though nobody called notify. This is not a bug in Rust; it is a documented reality of how condition variables map onto operating-system primitives, and the C, C++, and POSIX versions all share it. Re-checking the condition in a loop makes a spurious wake-up harmless -- the thread sees the condition is still false and simply goes back to sleep.
The contrast with the wrong tool is stark. Here is the tempting, wasteful busy-loop that a Condvar replaces:
use std::sync::{Arc, Mutex};
use std::thread;
fn main() {
let ready = Arc::new(Mutex::new(false));
let r2 = Arc::clone(&ready);
thread::spawn(move || { *r2.lock().unwrap() = true; });
// DON'T do this: a hot spin that pins a CPU core doing nothing useful
while !*ready.lock().unwrap() {
std::hint::spin_loop(); // burns cycles, lock/unlock churn, no sleeping
}
println!("finally ready"); // correct, but wasteful
}
It works, but it is exactly the CPU-burning anti-pattern we are here to avoid. The busy-loop grabs and releases the lock thousands of times a second, keeps a core spinning, and starves other work. The Condvar version sleeps at zero cost until the moment it matters.
Because the loop-around-wait pattern is so common, the standard library packages it as wait_while, which sleeps as long as the predicate holds and handles the spurious-wakeup loop for you:
use std::sync::{Arc, Mutex, Condvar};
use std::thread;
fn main() {
let pair = Arc::new((Mutex::new(0), Condvar::new()));
let p2 = Arc::clone(&pair);
thread::spawn(move || {
let (lock, cvar) = &*p2;
*lock.lock().unwrap() = 5;
cvar.notify_all();
});
let (lock, cvar) = &*pair;
// sleep while the count is still zero; guards against spurious wakeups automatically
let guard = cvar.wait_while(lock.lock().unwrap(), |count| *count == 0).unwrap();
println!("count reached {}", *guard); // count reached 5
}
Read wait_while(guard, predicate) as "keep sleeping while predicate is true". Its cousin, wait_until in some libraries, is the inverse phrasing; the standard library gives you wait_while, and once you internalise the direction it reads cleanly.
notify_one versus notify_all
You have two ways to wake sleepers, and choosing between them matters. notify_one wakes a single waiting thread; notify_all wakes every one of them. The rule of thumb: use notify_one when any single waiter can handle the event (one item was pushed, so one consumer should wake), and notify_all when the state change is relevant to everybody (a shutdown flag flipped, or a shared counter reached a value multiple waiters are checking). When in doubt, notify_all is the safe-but-slightly-wasteful choice -- it never leaves a thread wrongly asleep, at the cost of some threads waking, re-checking, and going back to sleep:
use std::sync::{Arc, Mutex, Condvar};
use std::thread;
fn main() {
let pair = Arc::new((Mutex::new(false), Condvar::new()));
let mut waiters = Vec::new();
for id in 0..3 {
let p = Arc::clone(&pair);
waiters.push(thread::spawn(move || {
let (lock, cvar) = &*p;
let _g = cvar.wait_while(lock.lock().unwrap(), |go| !*go).unwrap();
println!("waiter {id} woke and proceeds");
}));
}
let (lock, cvar) = &*pair;
*lock.lock().unwrap() = true;
cvar.notify_all(); // wake ALL three; notify_one would leave two asleep here
for w in waiters { w.join().unwrap(); }
}
Had we used notify_one here, only one of the three waiters would wake, and the program would hang on the two that never got their signal. The event ("everyone may go now") is relevant to all of them, so notify_all is correct.
A producer/consumer queue
Now the payoff. Put a Condvar next to a VecDeque and you have a blocking queue: the consumer sleeps when the queue is empty and wakes the instant the producer pushes. This is the beating heart of every channel implementation:
use std::sync::{Arc, Mutex, Condvar};
use std::collections::VecDeque;
use std::thread;
fn main() {
let shared = Arc::new((Mutex::new(VecDeque::new()), Condvar::new()));
let prod = Arc::clone(&shared);
let producer = thread::spawn(move || {
let (q, cvar) = &*prod;
for i in 0..5 {
q.lock().unwrap().push_back(i); // enqueue under the lock
cvar.notify_one(); // wake a sleeping consumer
}
});
let (q, cvar) = &*shared;
let mut received = Vec::new();
while received.len() < 5 {
let mut guard = q.lock().unwrap();
while guard.is_empty() {
guard = cvar.wait(guard).unwrap(); // sleep until an item shows up
}
received.push(guard.pop_front().unwrap());
}
producer.join().unwrap();
println!("consumed {} items", received.len()); // consumed 5 items
}
Trace the handshake once and it clicks. The consumer locks, sees the queue empty, and calls wait -- which releases the lock and sleeps. The producer can now lock, push, and notify_one. The consumer wakes, wait re-acquires the lock, the inner while guard.is_empty() re-checks (still guarding against spurious wakeups), finds an item, and pops it. This is, almost literally, what mpsc::channel from episode 63 does under the hood -- so you have now seen the machinery you have been standing on for the last ten episodes ;-)
The Barrier: a group rendezvous
A Condvar coordinates ad hoc signalling. A Barrier coordinates something more structured: a fixed group of threads that must all reach a synchronisation point before any of them proceeds. Every thread calls wait, and none returns from that call until the last one has arrived. It is the perfect tool for phased computations, where phase two must not begin until phase one is globally complete on every thread:
use std::sync::{Arc, Barrier};
use std::thread;
fn main() {
let barrier = Arc::new(Barrier::new(3));
thread::scope(|s| {
for id in 0..3 {
let b = Arc::clone(&barrier);
s.spawn(move || {
println!("worker {id} finished phase 1");
b.wait(); // no thread continues until all three reach here
println!("worker {id} starting phase 2");
});
}
});
}
Run this and you will see all three "finished phase 1" lines print (in some order) before any "starting phase 2" line appears. The barrier draws a hard line across time that no thread may cross early. Note the count you pass to Barrier::new(3) must match the number of threads that will call wait, or you deadlock -- a barrier of three with only two arrivals waits forever for a third that never comes.
A phased computation
The canonical use of a barrier is an iterative computation in rounds, where each round reads results the previous round produced. Every thread contributes its share, waits at the barrier so the shared state is fully settled, then safely reads the aggregate:
use std::sync::{Arc, Barrier, Mutex};
use std::thread;
fn main() {
let barrier = Arc::new(Barrier::new(4));
let sum = Arc::new(Mutex::new(0));
thread::scope(|s| {
for id in 1..=4 {
let (b, sum) = (Arc::clone(&barrier), Arc::clone(&sum));
s.spawn(move || {
*sum.lock().unwrap() += id; // phase 1: each thread contributes
b.wait(); // wait until ALL have contributed
let _total = *sum.lock().unwrap(); // phase 2 sees the complete sum
});
}
});
println!("total after phase 1: {}", *sum.lock().unwrap()); // 10
}
Without the barrier, a fast thread could race ahead into "phase 2" and read the sum while other threads were still adding to it -- a classic read-before-write bug. The barrier guarantees that by the time any thread reads _total, every contribution is already in. This is precisely the shape of parallel numeric solvers, cellular automata, and game-of-life-style simulations, where each generation must be complete before the next begins.
A count-down latch
A barrier is symmetric: everyone waits for everyone. Its useful asymmetric cousin is the count-down latch, where one thread waits for many to finish. You do not get one in std directly, but a counter, a Mutex, and a Condvar build one in a few lines. The main thread sleeps until N workers have each signalled completion:
use std::sync::{Arc, Mutex, Condvar};
use std::thread;
fn main() {
let latch = Arc::new((Mutex::new(3), Condvar::new())); // countdown starts at 3
for _ in 0..3 {
let l = Arc::clone(&latch);
thread::spawn(move || {
let (count, cvar) = &*l;
*count.lock().unwrap() -= 1; // this worker is done
cvar.notify_all(); // tell the waiter to re-check
});
}
let (count, cvar) = &*latch;
let guard = cvar.wait_while(count.lock().unwrap(), |c| *c > 0).unwrap();
println!("all workers done, count = {}", *guard); // all workers done, count = 0
}
Here notify_all is the right call even though only one thread is waiting: each worker fires it, and the single waiter re-checks *c > 0 on every wake, proceeding only when the count finally hits zero. This is the pattern behind Java's CountDownLatch and Go's sync.WaitGroup, which brings us neatly to the comparison.
How Python and Go would frame this
If you come from Python, the direct analogue of a Condvar is threading.Condition, and the loop-around-wait discipline is identical -- the Python docs themselves tell you to wait in a while, for the same spurious-wakeup reason:
import threading
cond = threading.Condition()
ready = False
def waiter():
with cond:
while not ready: # same loop-around-wait rule as Rust
cond.wait()
print("proceeding")
def signaller():
global ready
with cond:
ready = True
cond.notify() # like Rust's notify_one
t = threading.Thread(target=waiter)
t.start()
signaller()
t.join()
The with cond: block is Python's version of holding the mutex guard, cond.wait() releases and re-acquires exactly like Rust's wait, and notify/notify_all map one-to-one. The difference is that Python cannot stop you from forgetting to hold the lock, whereas Rust's wait literally consumes the MutexGuard and hands it back, so "you must hold the lock to wait" is enforced by the type system, not by documentation.
Go takes a different tack. It has sync.Cond, but idiomatic Go usually prefers channels for signalling and sync.WaitGroup for the count-down-latch job:
var wg sync.WaitGroup
for i := 0; i < 3; i++ {
wg.Add(1)
go func() {
defer wg.Done() // like decrementing our latch counter
// ... do work ...
}()
}
wg.Wait() // blocks until the counter hits zero, exactly like our latch
wg.Add(1) bumps the counter, wg.Done() decrements it, and wg.Wait() is our wait_while(count > 0) -- the whole latch we hand-built in Rust, baked into the standard library. Both framings are the same idea we implemented: a counter, atomic decrements, and a waiter that sleeps until it reaches zero. Seeing it in three languages should convince you the pattern is universal, even though each language dresses it differently.
When to reach for a Condvar
Here is the honest, practical guidance to close on. In real Rust you will not reach for a raw Condvar very often, and that is by design. Channels (from episode 63) cover most producer/consumer needs, and higher-level tools cover the rest -- so a hand-written Condvar is the exception, not the rule. But when you genuinely need to sleep until a specific shared state becomes true, and busy-looping is unacceptable (which it almost always is), the Condvar is exactly the right, efficient primitive, and nothing else in std does the job as directly.
The mental checklist is short. Do you need one thread to block until a condition another thread controls? Reach for (Mutex, Condvar) and always wait in a loop. Do you need a fixed group to rendezvous before proceeding in lockstep? Reach for Barrier. Do you need one thread to wait for many to finish? Build a count-down latch from a counter and a Condvar, or use a channel and count the messages. Get the loop-around-wait discipline into your fingers now, because you will see the exact same pattern again, in a very different guise, when we move from threads that block to tasks that yield -- coordinating not with sleeping OS threads, but with something that has to be initialised exactly once before anyone uses it.
Thanks for reading, and see you in the next one ;-)
We started by asking how last episode's recv() knew when to wake, and we end holding the answer: a Mutex to guard the state, a Condvar to sleep on it, and a stubborn while loop to shrug off spurious wakeups. From that one primitive we built a blocking queue, a barrier rendezvous, and a count-down latch -- the same three coordination shapes you will find in every threaded system ever written, in every language. Building them yourself is how you really understand them, so go do that and watch the threads fall into step. De groeten!
Exercises
- Use a
(Mutex<bool>, Condvar)pair so a worker thread sleeps until the main thread sets the flag totrueand callsnotify_one. Wait in a loop, not anif. - Use a
Barrier::new(3)so three threads each print "phase 1 done", then all print "phase 2 start" only after every thread has finished phase one. - Build a count-down latch where the main thread waits (via
wait_while) for four worker threads to each decrement a shared counter to zero before it prints a final message.