Skip to content

Rust Channels, Atomics and Worker Shutdown

Channels transfer messages; atomics update individual shared values; shutdown is a protocol built around those mechanisms. A stop flag does not wake a blocked receiver, and a successful send does not mean work has finished. This lesson makes those distinctions observable before moving to async code.

Prerequisites and outcome

Complete threads, Send, Sync, Arc and locks. You will drain a bounded multi-producer channel, join all workers, distinguish empty from disconnected, and explain why atomic ordering does not make a whole algorithm indivisible. Readings are synthetic Raspberry Pi millidegree fixtures; no device I/O or performance measurement is required.

Build the complete Rust channel example

cargo new pi_messages
cd pi_messages

Keep edition = "2024" in Cargo.toml. Replace src/main.rs with:

use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicI32, AtomicUsize, Ordering};
use std::sync::mpsc::{self, Receiver};
use std::thread;

#[derive(Debug, PartialEq, Eq)]
struct Collected {
    readings: Vec<(usize, Option<i32>)>,
    sent: usize,
}

fn collect_readings(
    batches: Vec<Vec<(usize, Option<i32>)>>,
    capacity: usize,
) -> Result<Collected, &'static str> {
    let (sender, receiver) = mpsc::sync_channel(capacity);
    let sent = Arc::new(AtomicUsize::new(0));
    let handles: Vec<_> = batches
        .into_iter()
        .map(|batch| {
            let sender = sender.clone();
            let counter = Arc::clone(&sent);
            thread::spawn(move || -> Result<(), &'static str> {
                for reading in batch {
                    sender.send(reading).map_err(|_| "receiver closed")?;
                    // This is a statistic, not a payload-publication signal.
                    counter.fetch_add(1, Ordering::Relaxed);
                }
                Ok(())
            })
        })
        .collect();
    drop(sender); // Do not retain a sender that prevents end-of-stream.

    // Drain before joining: bounded producers may be waiting for space.
    let mut readings: Vec<_> = receiver.into_iter().collect();
    let mut first_error = None;
    for handle in handles {
        let result = match handle.join() {
            Ok(result) => result,
            Err(_) => Err("producer panicked"),
        };
        if let Err(error) = result {
            first_error.get_or_insert(error);
        }
    }
    if let Some(error) = first_error {
        return Err(error);
    }
    readings.sort_by_key(|(index, _)| *index);
    Ok(Collected {
        readings,
        sent: sent.load(Ordering::Relaxed),
    })
}

enum Command {
    Record(i32),
    Stop,
}

fn command_worker(receiver: Receiver<Command>) -> Vec<i32> {
    let mut values = Vec::new();
    for command in receiver {
        match command {
            Command::Record(value) => values.push(value),
            Command::Stop => break,
        }
    }
    values
}

fn published_reading() -> i32 {
    let payload = Arc::new(AtomicI32::new(0));
    let ready = Arc::new(AtomicBool::new(false));
    let worker_payload = Arc::clone(&payload);
    let worker_ready = Arc::clone(&ready);
    let worker = thread::spawn(move || {
        worker_payload.store(46_700, Ordering::Relaxed);
        worker_ready.store(true, Ordering::Release);
    });
    // A tiny one-shot ordering demonstration, not a production waiting API.
    while !ready.load(Ordering::Acquire) {
        thread::yield_now();
    }
    let value = payload.load(Ordering::Relaxed);
    worker.join().expect("publication fixture panicked");
    value
}

fn main() {
    let batches = vec![vec![(0, Some(46_700)), (1, None)], vec![(2, Some(0))]];
    let report = collect_readings(batches, 1).expect("collection failed");
    println!("readings={:?}", report.readings);
    println!("sent={}", report.sent);

    let (sender, receiver) = mpsc::channel();
    let worker = thread::spawn(move || command_worker(receiver));
    sender.send(Command::Record(0)).unwrap();
    sender.send(Command::Record(-500)).unwrap();
    sender.send(Command::Stop).unwrap();
    drop(sender);
    println!("commands={:?}", worker.join().unwrap());
    println!("published={}", published_reading());
}

#[cfg(test)]
mod tests {
    use super::*;
    use std::sync::mpsc::{TryRecvError, TrySendError};

    #[test]
    fn bounded_collection_preserves_missing_zero_and_negative_values() {
        let batches = vec![vec![(2, Some(-500))], vec![(0, None), (1, Some(0))]];
        assert_eq!(
            collect_readings(batches, 1),
            Ok(Collected {
                readings: vec![(0, None), (1, Some(0)), (2, Some(-500))],
                sent: 3,
            })
        );
    }

    #[test]
    fn empty_and_rendezvous_collection_terminate() {
        assert_eq!(collect_readings(vec![], 0).unwrap().sent, 0);
        let report = collect_readings(vec![vec![], vec![(0, Some(0))]], 0).unwrap();
        assert_eq!(report.readings, vec![(0, Some(0))]);
        assert_eq!(report.sent, 1);
    }

    #[test]
    fn empty_is_not_disconnected_until_every_sender_is_dropped() {
        let (sender, receiver) = mpsc::channel::<i32>();
        let retained = sender.clone();
        drop(sender);
        assert_eq!(receiver.try_recv(), Err(TryRecvError::Empty));
        drop(retained);
        assert_eq!(receiver.try_recv(), Err(TryRecvError::Disconnected));
    }

    #[test]
    fn buffered_values_are_drained_before_disconnection() {
        let (sender, receiver) = mpsc::channel();
        sender.send(0).unwrap();
        drop(sender);
        assert_eq!(receiver.recv(), Ok(0));
        assert!(receiver.recv().is_err());
    }

    #[test]
    fn bounded_try_send_reports_full_and_returns_the_value() {
        let (sender, receiver) = mpsc::sync_channel(1);
        sender.send(String::from("first")).unwrap();
        let error = sender.try_send(String::from("second")).unwrap_err();
        assert!(matches!(error, TrySendError::Full(ref value) if value == "second"));
        assert_eq!(receiver.recv().unwrap(), "first");
        drop(receiver);
        let error = sender.send(String::from("third")).unwrap_err();
        assert_eq!(error.0, "third");
    }

    #[test]
    fn stop_and_disconnection_have_explicit_drain_policies() {
        let (sender, receiver) = mpsc::channel();
        sender.send(Command::Record(0)).unwrap();
        sender.send(Command::Stop).unwrap();
        sender.send(Command::Record(7)).unwrap();
        drop(sender);
        assert_eq!(command_worker(receiver), vec![0]);

        let (sender, receiver) = mpsc::channel();
        sender.send(Command::Record(7)).unwrap();
        drop(sender);
        assert_eq!(command_worker(receiver), vec![7]);
    }

    #[test]
    fn atomic_fetch_add_keeps_all_joined_increments() {
        let counter = Arc::new(AtomicUsize::new(0));
        let handles: Vec<_> = (0..4)
            .map(|_| {
                let counter = Arc::clone(&counter);
                thread::spawn(move || {
                    for _ in 0..250 {
                        counter.fetch_add(1, Ordering::Relaxed);
                    }
                })
            })
            .collect();
        for handle in handles {
            handle.join().unwrap();
        }
        assert_eq!(counter.load(Ordering::Relaxed), 1_000);
    }

    #[test]
    fn release_acquire_one_shot_publication_observes_payload() {
        assert_eq!(published_reading(), 46_700);
    }

    #[test]
    fn compare_exchange_returns_the_observed_previous_value() {
        let state = AtomicUsize::new(0);
        assert_eq!(
            state.compare_exchange(0, 1, Ordering::AcqRel, Ordering::Acquire),
            Ok(0)
        );
        assert_eq!(
            state.compare_exchange(0, 2, Ordering::AcqRel, Ordering::Acquire),
            Err(1)
        );
        assert_eq!(state.load(Ordering::Relaxed), 1);
    }
}
1
2
3
4
cargo check
cargo test
cargo fmt --check
cargo run --quiet

Expected binary output:

1
2
3
4
readings=[(0, Some(46700)), (1, None), (2, Some(0))]
sent=3
commands=[0, -500]
published=46700

Nine tests cover message accounting, empty streams, rendezvous, sender lifetime, buffered draining, backpressure, stop policy, atomic increments, publication and compare-exchange. Message arrival from different producers is not predetermined; the fixture sorts unique indices solely to give reproducible displayed output.

Channels transfer ownership and distinguish two kinds of waiting

The standard mpsc channel has cloneable senders and one receiver. Sending an owned String transfers it; sending an Arc transfers a shared-ownership handle, not a deep copy of its payload. The channel APIs and thread synchronization are standard-library mechanisms, while the move and borrowing constraints remain language rules. See mpsc.

channel has an unbounded queue in its API model; it is not a promise of unlimited physical RAM. sync_channel uses a configured capacity and a send can wait for space. Capacity zero is a rendezvous: sending waits for a corresponding receive. This bounds queued messages, not necessarily all allocations held by workers. See sync_channel.

recv waits for a value or disconnection. try_recv does not wait and distinguishes Empty from Disconnected. Dropping all senders permits the receiver to finish after queued messages are drained; a forgotten sender clone keeps an otherwise empty receive loop waiting. recv_timeout adds a timeout outcome, not proof that producers are dead. See Receiver.

Sending can fail when the receiver is gone. SendError retains the unsent value, so a caller can choose recovery rather than silently lose it. A successful send only means the message was accepted by the channel. A separate reply, acknowledgement or joined worker result is needed when completion matters.

Drain bounded producers before joining them

Each producer owns a batch, sends indexed records and increments a statistic after a successful send. Main drops its unused sender, consumes messages, and then joins every producer, including after an observed worker error. Joining before receiving could create a wait cycle: a full queue blocks a producer, while main waits for that producer to finish.

The final sent count is read after joins, so the workers have finished their updates. During collection it could lag behind received messages because the increment follows send. It is not a completion signal. Sorting requires unique fixture indices and does not establish meaningful chronological ordering for unrelated producers.

For a sustained service, also define thread-creation failure handling, bounded input production, duplicate identifiers and what partial results mean after worker failure. This finite fixture is not a general task pool or a production shutdown framework.

Shutdown is a protocol, not merely a Boolean

command_worker accepts Record or Stop. This policy processes records preceding Stop from that sender, stops at Stop, and abandons later queued records. Without Stop, disconnection drains the queue. The two cases are tested separately; neither is automatically the right policy for every application.

A protocol should answer: who stops accepting input, which queued work is completed, how blocked operations become able to finish, and who joins each worker. With multiple producers, one producer's Stop does not prove all other producers have finished. Multiple consumers need a different work-distribution design: std mpsc Receiver is not Clone or Sync.

An atomic cancellation flag is cooperative. Workers must check it at defined points; it does not kill a thread, roll back work or wake recv. Use an appropriate control message, closure of all input handles, or an explicit timeout/wakeup strategy for blocked workers. Similarly, a blocked bounded send cannot be interrupted merely by changing an unrelated Boolean. Do not use sleeps as a shutdown guarantee.

Atomics make operations indivisible, not entire algorithms

AtomicUsize::fetch_add performs one read-modify-write operation. A separate load followed by store can lose another worker's update even with SeqCst, because the pair is still two operations. For counters, use a suitable atomic RMW; for coupled fields and invariants, a mutex or message owner is often clearer. See AtomicUsize.

Atomic types have target-dependent availability and distinct APIs; they are not generic wrappers for arbitrary Rust values. They do not by themselves guarantee a multi-step algorithm's lock-free progress. Our Pi's supported integer atomics are used here; do not infer that every target offers every width. See the atomic module.

compare_exchange conditionally replaces one atomic value. Ok contains the old value that matched the expectation; Err contains the observed value that did not. Its success and failure orderings describe different paths, and compare_exchange_weak may fail spuriously. A retry loop must reconsider its assumptions, not blindly reuse stale state; logical protocols can still suffer from value changes that return to the same representation (the ABA problem). The atomic operation APIs describe the supported orderings and return values.

Choose ordering from the synchronization relationship

Relaxed preserves atomicity without publishing unrelated memory. It suits our statistic because it does not guard a payload and main subsequently joins producers. Release on a store and Acquire on a load establish synchronization when that load observes the relevant released value. AcqRel combines those directions for an RMW; SeqCst additionally orders sequentially consistent atomic operations globally. See Ordering.

In published_reading, the payload's store precedes a Release store to ready. The reader waits until its Acquire load sees true, then reads the payload before joining. Both fields are atomic, so even an incorrectly weakened ordering variant would not create a non-atomic data race; it would fail to establish the intended cross-field guarantee. The flag is one-shot: a reusable publication protocol needs additional reasoning about repeated updates and ownership.

Acquire is a load-side ordering and Release is a store-side ordering; invalid combinations can panic and constant invalid choices may be rejected by a compiler lint. compare_exchange's failure path only loads and cannot use Release or AcqRel. Read the operation's API rather than attaching the strongest-looking name arbitrarily. Our fixture uses AtomicBool for the flag.

The yield loop is a tiny ordering demonstration, not a recommended waiting strategy. It has no timeout or recovery path if publication never occurs, and scheduling is not guaranteed by yield_now. Prefer channels or other waiting primitives for application-level events. Passing tests on a Pi does not prove a weak-memory algorithm; rely on the API's specified synchronization relations, not whether a bad variant happens to fail during a run.

Deliberately failing: send consumes an owned message

In a separate scratch project, replace src/main.rs with:

1
2
3
4
5
6
7
8
9
use std::sync::mpsc;

fn main() {
    let (sender, receiver) = mpsc::channel();
    let label = String::from("Pi 4B");
    sender.send(label).unwrap();
    println!("{label}");
    drop(receiver);
}

cargo check reports E0382. The message moved into send. Read it from the receiver, clone deliberately before sending if two owned copies are required, or choose a different ownership model.

Exercises and troubleshooting

  1. Repair the scratch program by printing receiver.recv().unwrap() instead of label. Expect Pi 4B. Explain why send may succeed before the receiver has consumed the message.
  2. Attempt receiver.clone() in a scratch program: expect E0599. Explain why cloneable senders do not imply cloneable receivers.
  3. Change the normal bounded capacity from 1 to 0 or 4. The displayed report should be unchanged. Discuss which send operations may wait, without treating a timing observation as an API guarantee.
  4. Use try_recv with one retained sender clone, then drop that clone. Expect Empty followed by Disconnected; this safely demonstrates a hang's cause without running an infinite receive loop.
  5. Replace an atomic increment with load followed by store. Two workers can both read zero before either stores one. Even SeqCst allows the final one instead of two; a barrier-controlled negative test can demonstrate this interleaving without relying on luck. Repair with fetch_add.
  6. Explain why setting a stop flag cannot wake a worker blocked in recv. Specify an actual wakeup/closure mechanism and decide whether pending work is drained or abandoned.
  7. Change the one-shot ready operations to Relaxed. Do not claim that an observed correct output proves publication: state which synchronization guarantee has been removed. Do not introduce unsafely shared non-atomic data to try to manufacture a failure.

Verification and next step

On October 10, 2026, this lesson was verified on the authorised Raspberry Pi 4B with 64-bit user space, kernel 6.18.50+rpt-rpi-v8, Rust and Cargo 1.99.0, and edition 2024. Cargo check, nine debug/release tests, formatting, debug/release outputs and five further debug output comparisons passed. Capacities zero and four preserved the report, and receive-after-send plus atomic-RMW repairs passed. Use after send and receiver cloning produced E0382 and E0599. A barrier-controlled load/store counter compiled but failed its expected-two assertion; replacing its store with fetch_add passed. These finite runs demonstrate the chosen cases, not every memory interleaving, production cancellation or a speed improvement. The publication guarantee is justified by the Release/Acquire API contract, not by repeated observed output.

Next, study async, await and Future to distinguish suspended computations from threads and runtime scheduling.

Previous: threads and locks · Course overview

Donate