Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 4 additions & 4 deletions futures-channel/src/mpsc/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -521,7 +521,7 @@ impl<T> BoundedSenderInner<T> {
/// if there was an error.
fn try_send(&mut self, msg: T) -> Result<(), TrySendError<T>> {
// If the sender is currently blocked, reject the message
if !self.poll_unparked(None).is_ready() {
if self.poll_unparked(None).is_pending() {
return Err(TrySendError { err: SendError { kind: SendErrorKind::Full }, val: msg });
}

Expand All @@ -545,9 +545,9 @@ impl<T> BoundedSenderInner<T> {
// receiver is dropped.
let park_self = match self.inc_num_messages() {
Some(num_messages) => {
// Block if the current number of pending messages has exceeded
// the configured buffer size
num_messages > self.inner.buffer
// Block next messages if the current number of pending messages
// has reached the configured capacity (buffer + num_senders)
num_messages >= self.inner.buffer + self.inner.num_senders.load(SeqCst)
}
None => {
return Err(TrySendError {
Expand Down
12 changes: 4 additions & 8 deletions futures-channel/src/mpsc/sink_impl.rs
Original file line number Diff line number Diff line change
Expand Up @@ -14,14 +14,10 @@ impl<T> Sink<T> for Sender<T> {
(*self).start_send(msg)
}

fn poll_flush(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
match (*self).poll_ready(cx) {
Poll::Ready(Err(ref e)) if e.is_disconnected() => {
// If the receiver disconnected, we consider the sink to be flushed.
Poll::Ready(Ok(()))
}
x => x,
}
fn poll_flush(self: Pin<&mut Self>, _: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
// MPSC channels have no private internal buffering in the sender; start_send
// immediately enqueues messages to receiver, so there is nothing to flush.
Poll::Ready(Ok(()))
}

fn poll_close(mut self: Pin<&mut Self>, _: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
Expand Down
111 changes: 77 additions & 34 deletions futures-channel/tests/mpsc.rs
Original file line number Diff line number Diff line change
@@ -1,12 +1,11 @@
use futures::channel::{mpsc, oneshot};
use futures::executor::{block_on, block_on_stream};
use futures::future::{poll_fn, FutureExt};
use futures::sink::{Sink, SinkExt};
use futures::sink::SinkExt;
use futures::stream::{Stream, StreamExt};
use futures::task::{Context, Poll};
use futures_channel::mpsc::{RecvError, TryRecvError};
use futures_test::task::{new_count_waker, noop_context};
use std::pin::pin;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use std::thread;
Expand All @@ -27,40 +26,80 @@ fn send_recv() {
}

#[test]
fn send_recv_no_buffer() {
// Run on a task context
block_on(poll_fn(move |cx| {
let (tx, rx) = mpsc::channel::<i32>(0);
let mut tx = pin!(tx);
let mut rx = pin!(rx);
fn send_recv_buffer_zero() {
// Test that channel capacity = buffer + num_senders
// With buffer=0 and 1 sender, capacity should be 1
let (mut tx, rx) = mpsc::channel::<i32>(0);
let mut rx = block_on_stream(rx);

// First send should succeed (capacity = 0 + 1 sender = 1)
tx.try_send(1).unwrap();

assert!(tx.as_mut().poll_flush(cx).is_ready());
assert!(tx.as_mut().poll_ready(cx).is_ready());
// Second send should fail - channel is full
assert!(tx.try_send(2).is_err());

// Send first message
assert!(tx.as_mut().start_send(1).is_ok());
assert!(tx.as_mut().poll_ready(cx).is_pending());
// Clone the sender - now capacity should be 2 (0 + 2 senders)
let mut tx2 = tx.clone();

// poll_ready said Pending, so no room in buffer, therefore new sends
// should get rejected with is_full.
assert!(tx.as_mut().start_send(0).unwrap_err().is_full());
assert!(tx.as_mut().poll_ready(cx).is_pending());
// Send on cloned sender should succeed (it's not parked, and capacity increased)
tx2.try_send(2).unwrap();

// Take the value
assert_eq!(rx.as_mut().poll_next(cx), Poll::Ready(Some(1)));
assert!(tx.as_mut().poll_ready(cx).is_ready());
// Third send should fail - channel is full again
assert!(tx2.try_send(3).is_err());

// Send second message
assert!(tx.as_mut().poll_ready(cx).is_ready());
assert!(tx.as_mut().start_send(2).is_ok());
assert!(tx.as_mut().poll_ready(cx).is_pending());
// Receive messages in order
assert_eq!(rx.next(), Some(1));
assert_eq!(rx.next(), Some(2));

// Take the value
assert_eq!(rx.as_mut().poll_next(cx), Poll::Ready(Some(2)));
assert!(tx.as_mut().poll_ready(cx).is_ready());
// Now the pending task should be able to complete
tx2.try_send(3).unwrap();
assert_eq!(rx.next(), Some(3));

Poll::Ready(())
}));
// Drop both senders
drop(tx);
drop(tx2);

// Receiver should indicate stream ended
assert_eq!(rx.next(), None);
}

#[test]
fn send_recv_buffer_zero_async() {
// Same as send_recv_buffer_zero but exercising the async send path (Sink trait)
// Tests that capacity = buffer + num_senders works correctly through poll_ready/start_send/poll_flush
let (mut tx, rx) = mpsc::channel::<i32>(0);
let mut rx = block_on_stream(rx);
let mut cx = noop_context();

// First send should succeed (capacity = 0 + 1 sender = 1)
block_on(tx.send(1)).unwrap();

// Second send should be pending - channel is full with one sender
assert_eq!(tx.send(2).poll_unpin(&mut cx), Poll::Pending);

// Clone the sender - now capacity should be 2 (0 + 2 senders)
let mut tx2 = tx.clone();

// Send on cloned sender should succeed (capacity increased)
block_on(tx2.send(2)).unwrap();

// Now try to send again on the original sender - should be pending (channel full again)
assert_eq!(tx.send(3).poll_unpin(&mut cx), Poll::Pending);

// Receive messages in order
assert_eq!(rx.next(), Some(1));
assert_eq!(rx.next(), Some(2));

// Now the pending task should be able to complete
assert_eq!(tx.send(3).poll_unpin(&mut cx), Poll::Ready(Ok(())));
assert_eq!(rx.next(), Some(3));

// Drop both senders
drop(tx);
drop(tx2);

// Receiver should indicate stream ended
assert_eq!(rx.next(), None);
}

#[test]
Expand Down Expand Up @@ -332,7 +371,7 @@ fn stress_drop_sender() {
const ITER: usize = if cfg!(miri) { 100 } else { 10000 };

fn list() -> impl Stream<Item = i32> {
let (tx, rx) = mpsc::channel(1);
let (tx, rx) = mpsc::channel(0);
thread::spawn(move || {
block_on(send_one_two_three(tx));
});
Expand Down Expand Up @@ -432,7 +471,7 @@ fn stress_poll_ready() {
#[test]
fn test_bounded_recv() {
let (dropped_tx, dropped_rx) = oneshot::channel();
let (tx, mut rx) = mpsc::channel(1);
let (tx, mut rx) = mpsc::channel(0);
thread::spawn(move || {
block_on(async move {
send_one_two_three(tx).await;
Expand Down Expand Up @@ -475,7 +514,7 @@ fn test_unbounded_recv() {
#[test]
fn try_send_1() {
const N: usize = if cfg!(miri) { 100 } else { 3000 };
let (mut tx, rx) = mpsc::channel(0);
let (mut tx, rx) = mpsc::channel(1);

let t = thread::spawn(move || {
for i in 0..N {
Expand Down Expand Up @@ -639,7 +678,7 @@ fn send_backpressure() {
let (waker, counter) = new_count_waker();
let mut cx = Context::from_waker(&waker);

let (mut tx, mut rx) = mpsc::channel(1);
let (mut tx, mut rx) = mpsc::channel(0);
block_on(tx.send(1)).unwrap();

let mut task = tx.send(2);
Expand All @@ -660,11 +699,15 @@ fn send_backpressure_multi_senders() {
let (waker, counter) = new_count_waker();
let mut cx = Context::from_waker(&waker);

let (mut tx1, mut rx) = mpsc::channel(1);
let (mut tx1, mut rx) = mpsc::channel(0);
let mut tx2 = tx1.clone();
block_on(tx1.send(1)).unwrap();

let mut task = tx2.send(2);
assert_eq!(task.poll_unpin(&mut cx), Poll::Ready(Ok(())));
assert_eq!(counter, 0);

let mut task = tx2.send(3);
assert_eq!(task.poll_unpin(&mut cx), Poll::Pending);
assert_eq!(counter, 0);

Expand Down
Loading