Skip to content

Instantly share code, notes, and snippets.

@AnthoniG
Forked from jonhoo/mpsc.rs
Created September 7, 2026 18:36
Show Gist options
  • Select an option

  • Save AnthoniG/e6f65c5cb98f55d2bd906cb5eeccb28f to your computer and use it in GitHub Desktop.

Select an option

Save AnthoniG/e6f65c5cb98f55d2bd906cb5eeccb28f to your computer and use it in GitHub Desktop.
use std::collections::VecDeque;
use std::sync::{Arc, Condvar, Mutex};
// Flavors:
// - Synchronous channels: Channel where send() can block. Limited capacity.
// - Mutex + Condvar + VecDeque
// - Atomic VecDeque (atomic queue) + thread::park + thread::Thread::notify
// - Asynchronous channels: Channel where send() cannot block. Unbounded.
// - Mutex + Condvar + VecDeque
// - Mutex + Condvar + LinkedList
// - Atomic linked list, linked list of T
// - Atomic block linked list, linked list of atomic VecDeque<T>
// - Rendezvous channels: Synchronous with capacity = 0. Used for thread synchronization.
// - Oneshot channels: Any capacity. In practice, only one call to send().
// async/await
pub struct Sender<T> {
shared: Arc<Shared<T>>,
}
impl<T> Clone for Sender<T> {
fn clone(&self) -> Self {
let mut inner = self.shared.inner.lock().unwrap();
inner.senders += 1;
drop(inner);
Sender {
shared: Arc::clone(&self.shared),
}
}
}
impl<T> Drop for Sender<T> {
fn drop(&mut self) {
let mut inner = self.shared.inner.lock().unwrap();
inner.senders -= 1;
let was_last = inner.senders == 0;
drop(inner);
if was_last {
self.shared.available.notify_one();
}
}
}
impl<T> Sender<T> {
pub fn send(&mut self, t: T) {
let mut inner = self.shared.inner.lock().unwrap();
inner.queue.push_back(t);
drop(inner);
self.shared.available.notify_one();
}
}
pub struct Receiver<T> {
shared: Arc<Shared<T>>,
buffer: VecDeque<T>,
}
impl<T> Receiver<T> {
pub fn recv(&mut self) -> Option<T> {
if let Some(t) = self.buffer.pop_front() {
return Some(t);
}
let mut inner = self.shared.inner.lock().unwrap();
loop {
match inner.queue.pop_front() {
Some(t) => {
std::mem::swap(&mut self.buffer, &mut inner.queue);
return Some(t);
}
None if inner.senders == 0 => return None,
None => {
inner = self.shared.available.wait(inner).unwrap();
}
}
}
}
}
impl<T> Iterator for Receiver<T> {
type Item = T;
fn next(&mut self) -> Option<Self::Item> {
self.recv()
}
}
struct Inner<T> {
queue: VecDeque<T>,
senders: usize,
}
struct Shared<T> {
inner: Mutex<Inner<T>>,
available: Condvar,
}
pub fn channel<T>() -> (Sender<T>, Receiver<T>) {
let inner = Inner {
queue: VecDeque::default(),
senders: 1,
};
let shared = Shared {
inner: Mutex::new(inner),
available: Condvar::new(),
};
let shared = Arc::new(shared);
(
Sender {
shared: shared.clone(),
},
Receiver {
shared: shared.clone(),
},
)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn ping_pong() {
let (mut tx, mut rx) = channel();
tx.send(42);
assert_eq!(rx.recv(), Some(42));
}
#[test]
fn closed_tx() {
let (tx, mut rx) = channel::<()>();
drop(tx);
assert_eq!(rx.recv(), None);
}
#[test]
fn closed_rx() {
let (mut tx, rx) = channel();
drop(rx);
tx.send(42);
}
}
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment