mirror of
https://github.com/tokio-rs/tokio.git
synced 2026-08-21 00:00:10 +02:00
task: improve the example of poll_proceed (#7586)
Signed-off-by: ADD-SP <[email protected]>
This commit is contained in:
+22
-58
@@ -306,74 +306,38 @@ cfg_coop! {
|
||||
///
|
||||
/// # Examples
|
||||
///
|
||||
/// This example shows a simple countdown latch that uses [`poll_proceed`] to participate in
|
||||
/// cooperative scheduling.
|
||||
/// This example wraps the `futures::channel::mpsc::UnboundedReceiver` to
|
||||
/// cooperate with the Tokio scheduler. Each time a value is received, task budget
|
||||
/// is consumed. If no budget is available, the task yields to the scheduler.
|
||||
///
|
||||
/// ```
|
||||
/// use std::future::{Future};
|
||||
/// use std::pin::Pin;
|
||||
/// use std::task::{ready, Context, Poll, Waker};
|
||||
/// use std::task::{ready, Context, Poll};
|
||||
/// use tokio::task::coop;
|
||||
/// use futures::stream::{Stream, StreamExt};
|
||||
/// use futures::channel::mpsc::UnboundedReceiver;
|
||||
///
|
||||
/// struct CountdownLatch<T> {
|
||||
/// counter: usize,
|
||||
/// value: Option<T>,
|
||||
/// waker: Option<Waker>
|
||||
/// struct CoopUnboundedReceiver<T> {
|
||||
/// receiver: UnboundedReceiver<T>,
|
||||
/// }
|
||||
///
|
||||
/// impl<T> CountdownLatch<T> {
|
||||
/// fn new(value: T, count: usize) -> Self {
|
||||
/// CountdownLatch {
|
||||
/// counter: count,
|
||||
/// value: Some(value),
|
||||
/// waker: None
|
||||
/// }
|
||||
/// }
|
||||
/// fn count_down(&mut self) {
|
||||
/// if self.counter <= 0 {
|
||||
/// return;
|
||||
/// }
|
||||
///
|
||||
/// self.counter -= 1;
|
||||
/// if self.counter == 0 {
|
||||
/// if let Some(w) = self.waker.take() {
|
||||
/// w.wake();
|
||||
/// }
|
||||
/// }
|
||||
/// }
|
||||
/// }
|
||||
///
|
||||
/// impl<T> Future for CountdownLatch<T> {
|
||||
/// type Output = T;
|
||||
///
|
||||
/// fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
|
||||
/// // `poll_proceed` checks with the runtime if this task is still allowed to proceed
|
||||
/// // with performing work.
|
||||
/// // If not, `Pending` is returned and `ready!` ensures this function returns.
|
||||
/// // If we are allowed to proceed, coop now represents the budget consumption
|
||||
/// impl<T> Stream for CoopUnboundedReceiver<T> {
|
||||
/// type Item = T;
|
||||
/// fn poll_next(
|
||||
/// mut self: Pin<&mut Self>,
|
||||
/// cx: &mut Context<'_>
|
||||
/// ) -> Poll<Option<T>> {
|
||||
/// let coop = ready!(coop::poll_proceed(cx));
|
||||
///
|
||||
/// // Get a mutable reference to the CountdownLatch
|
||||
/// let this = Pin::get_mut(self);
|
||||
///
|
||||
/// // Next we check if the latch is ready to release its value
|
||||
/// if this.counter == 0 {
|
||||
/// let t = this.value.take();
|
||||
/// // The latch made progress so call `made_progress` to ensure the budget
|
||||
/// // is not reverted.
|
||||
/// coop.made_progress();
|
||||
/// Poll::Ready(t.unwrap())
|
||||
/// } else {
|
||||
/// // If the latch is not ready so return pending and simply drop `coop`.
|
||||
/// // This will restore the budget making it available again to perform any
|
||||
/// // other work.
|
||||
/// this.waker = Some(cx.waker().clone());
|
||||
/// Poll::Pending
|
||||
/// }
|
||||
/// match self.receiver.poll_next_unpin(cx) {
|
||||
/// Poll::Ready(v) => {
|
||||
/// // We received a value, so consume budget.
|
||||
/// coop.made_progress();
|
||||
/// Poll::Ready(v)
|
||||
/// }
|
||||
/// Poll::Pending => Poll::Pending,
|
||||
/// }
|
||||
/// }
|
||||
/// }
|
||||
///
|
||||
/// impl<T> Unpin for CountdownLatch<T> {}
|
||||
/// ```
|
||||
#[inline]
|
||||
pub fn poll_proceed(cx: &mut Context<'_>) -> Poll<RestoreOnPending> {
|
||||
|
||||
Reference in New Issue
Block a user