mirror of
https://github.com/tokio-rs/tokio.git
synced 2026-09-09 00:00:08 +02:00
async-await: track nightly changes (#661)
The `tokio-async-await` crate is no longer a facade. Instead, the `tokio` crate provides a feature flag to enable async/await support.
This commit is contained in:
@@ -1,112 +0,0 @@
|
||||
use futures::{
|
||||
Future as Future01,
|
||||
Poll as Poll01,
|
||||
};
|
||||
use futures_core::{Future as Future03};
|
||||
|
||||
use std::pin::PinBox;
|
||||
use std::future::FutureObj;
|
||||
use std::ptr::NonNull;
|
||||
use std::task::{
|
||||
Context,
|
||||
Spawn,
|
||||
UnsafeWake,
|
||||
LocalWaker,
|
||||
Poll as Poll03,
|
||||
Waker,
|
||||
SpawnObjError,
|
||||
};
|
||||
|
||||
/// Convert an 0.3 `Future` to an 0.1 `Future`.
|
||||
#[derive(Debug)]
|
||||
pub struct Compat<T>(PinBox<T>);
|
||||
|
||||
impl<T> Compat<T> {
|
||||
pub fn new(data: T) -> Compat<T> {
|
||||
Compat(PinBox::new(data))
|
||||
}
|
||||
}
|
||||
|
||||
/// Convert a value into one that can be used with `await!`.
|
||||
pub trait IntoAwaitable {
|
||||
type Awaitable;
|
||||
|
||||
fn into_awaitable(self) -> Self::Awaitable;
|
||||
}
|
||||
|
||||
impl<T> IntoAwaitable for T
|
||||
where T: Future03,
|
||||
{
|
||||
type Awaitable = Self;
|
||||
|
||||
fn into_awaitable(self) -> Self {
|
||||
self
|
||||
}
|
||||
}
|
||||
|
||||
impl<T, Item, Error> Future01 for Compat<T>
|
||||
where T: Future03<Output = Result<Item, Error>>,
|
||||
{
|
||||
type Item = Item;
|
||||
type Error = Error;
|
||||
|
||||
fn poll(&mut self) -> Poll01<Item, Error> {
|
||||
use futures::Async::*;
|
||||
|
||||
let local_waker = noop_local_waker();
|
||||
let mut executor = NoopExecutor;
|
||||
|
||||
let mut cx = Context::new(&local_waker, &mut executor);
|
||||
|
||||
let res = self.0.as_pin_mut().poll(&mut cx);
|
||||
|
||||
match res {
|
||||
Poll03::Ready(Ok(val)) => Ok(Ready(val)),
|
||||
Poll03::Ready(Err(err)) => Err(err),
|
||||
Poll03::Pending => Ok(NotReady),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// ===== NoopWaker =====
|
||||
|
||||
struct NoopWaker;
|
||||
|
||||
fn noop_local_waker() -> LocalWaker {
|
||||
let w: NonNull<NoopWaker> = NonNull::dangling();
|
||||
unsafe { LocalWaker::new(w) }
|
||||
}
|
||||
|
||||
fn noop_waker() -> Waker {
|
||||
let w: NonNull<NoopWaker> = NonNull::dangling();
|
||||
unsafe { Waker::new(w) }
|
||||
}
|
||||
|
||||
unsafe impl UnsafeWake for NoopWaker {
|
||||
unsafe fn clone_raw(&self) -> Waker {
|
||||
noop_waker()
|
||||
}
|
||||
|
||||
unsafe fn drop_raw(&self) {
|
||||
}
|
||||
|
||||
unsafe fn wake(&self) {
|
||||
panic!("NoopWake cannot wake");
|
||||
}
|
||||
}
|
||||
|
||||
// ===== NoopExecutor =====
|
||||
|
||||
struct NoopExecutor;
|
||||
|
||||
impl Spawn for NoopExecutor {
|
||||
fn spawn_obj(&mut self, future: FutureObj<'static, ()>) -> Result<(), SpawnObjError> {
|
||||
use std::task::SpawnErrorKind;
|
||||
|
||||
// NoopExecutor cannot execute
|
||||
Err(SpawnObjError {
|
||||
kind: SpawnErrorKind::shutdown(),
|
||||
future,
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -1,69 +0,0 @@
|
||||
|
||||
use futures::{Future, Async};
|
||||
use futures_core::future::Future as Future03;
|
||||
use futures_core::task::Poll as Poll03;
|
||||
|
||||
use std::marker::Unpin;
|
||||
use std::pin::PinMut;
|
||||
use std::task::Context;
|
||||
|
||||
/// Converts an 0.1 `Future` into an 0.3 `Future`.
|
||||
#[derive(Debug)]
|
||||
pub struct Compat<T>(T);
|
||||
|
||||
pub(crate) fn convert_poll<T, E>(poll: Result<Async<T>, E>) -> Poll03<Result<T, E>> {
|
||||
use futures::Async::{Ready, NotReady};
|
||||
|
||||
match poll {
|
||||
Ok(Ready(val)) => Poll03::Ready(Ok(val)),
|
||||
Ok(NotReady) => Poll03::Pending,
|
||||
Err(err) => Poll03::Ready(Err(err)),
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) fn convert_poll_stream<T, E>(
|
||||
poll: Result<Async<Option<T>>, E>) -> Poll03<Option<Result<T, E>>>
|
||||
{
|
||||
use futures::Async::{Ready, NotReady};
|
||||
|
||||
match poll {
|
||||
Ok(Ready(Some(val))) => Poll03::Ready(Some(Ok(val))),
|
||||
Ok(Ready(None)) => Poll03::Ready(None),
|
||||
Ok(NotReady) => Poll03::Pending,
|
||||
Err(err) => Poll03::Ready(Some(Err(err))),
|
||||
}
|
||||
}
|
||||
|
||||
/// Convert a value into one that can be used with `await!`.
|
||||
pub trait IntoAwaitable {
|
||||
type Awaitable;
|
||||
|
||||
/// Convert `self` into a value that can be used with `await!`.
|
||||
fn into_awaitable(self) -> Self::Awaitable;
|
||||
}
|
||||
|
||||
impl<T: Future + Unpin> IntoAwaitable for T {
|
||||
type Awaitable = Compat<T>;
|
||||
|
||||
fn into_awaitable(self) -> Self::Awaitable {
|
||||
Compat(self)
|
||||
}
|
||||
}
|
||||
|
||||
impl<T> Future03 for Compat<T>
|
||||
where T: Future + Unpin
|
||||
{
|
||||
type Output = Result<T::Item, T::Error>;
|
||||
|
||||
fn poll(self: PinMut<Self>, _cx: &mut Context) -> Poll03<Self::Output> {
|
||||
use futures::Async::{Ready, NotReady};
|
||||
|
||||
// TODO: wire in cx
|
||||
|
||||
match PinMut::get_mut(self).0.poll() {
|
||||
Ok(Ready(val)) => Poll03::Ready(Ok(val)),
|
||||
Ok(NotReady) => Poll03::Pending,
|
||||
Err(e) => Poll03::Ready(Err(e)),
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,8 +0,0 @@
|
||||
//! Utilities for working with `async` / `await`.
|
||||
|
||||
#[macro_use]
|
||||
mod await;
|
||||
pub mod compat;
|
||||
pub mod io;
|
||||
pub mod sink;
|
||||
pub mod stream;
|
||||
@@ -3,8 +3,10 @@
|
||||
macro_rules! await {
|
||||
($e:expr) => {{
|
||||
use $crate::std_await;
|
||||
use $crate::async_await::compat::forward::IntoAwaitable as IntoAwaitableForward;
|
||||
use $crate::async_await::compat::backward::IntoAwaitable as IntoAwaitableBackward;
|
||||
#[allow(unused_imports)]
|
||||
use $crate::compat::forward::IntoAwaitable as IntoAwaitableForward;
|
||||
#[allow(unused_imports)]
|
||||
use $crate::compat::backward::IntoAwaitable as IntoAwaitableBackward;
|
||||
|
||||
#[allow(unused_mut)]
|
||||
let mut e = $e;
|
||||
@@ -0,0 +1,89 @@
|
||||
use futures::{Future, Poll};
|
||||
|
||||
use std::pin::Pin;
|
||||
use std::future::{
|
||||
Future as StdFuture,
|
||||
};
|
||||
use std::ptr::NonNull;
|
||||
use std::task::{
|
||||
LocalWaker,
|
||||
Poll as StdPoll,
|
||||
UnsafeWake,
|
||||
Waker,
|
||||
};
|
||||
|
||||
/// Convert an 0.3 `Future` to an 0.1 `Future`.
|
||||
#[derive(Debug)]
|
||||
pub struct Compat<T>(Pin<Box<T>>);
|
||||
|
||||
impl<T> Compat<T> {
|
||||
/// Create a new `Compat` backed by `future`.
|
||||
pub fn new(future: T) -> Compat<T> {
|
||||
Compat(Box::pinned(future))
|
||||
}
|
||||
}
|
||||
|
||||
/// Convert a value into one that can be used with `await!`.
|
||||
pub trait IntoAwaitable {
|
||||
type Awaitable;
|
||||
|
||||
fn into_awaitable(self) -> Self::Awaitable;
|
||||
}
|
||||
|
||||
impl<T> IntoAwaitable for T
|
||||
where T: StdFuture,
|
||||
{
|
||||
type Awaitable = Self;
|
||||
|
||||
fn into_awaitable(self) -> Self {
|
||||
self
|
||||
}
|
||||
}
|
||||
|
||||
impl<T, Item, Error> Future for Compat<T>
|
||||
where T: StdFuture<Output = Result<Item, Error>>,
|
||||
{
|
||||
type Item = Item;
|
||||
type Error = Error;
|
||||
|
||||
fn poll(&mut self) -> Poll<Item, Error> {
|
||||
use futures::Async::*;
|
||||
|
||||
let local_waker = noop_local_waker();
|
||||
|
||||
let res = self.0.as_mut().poll(&local_waker);
|
||||
|
||||
match res {
|
||||
StdPoll::Ready(Ok(val)) => Ok(Ready(val)),
|
||||
StdPoll::Ready(Err(err)) => Err(err),
|
||||
StdPoll::Pending => Ok(NotReady),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// ===== NoopWaker =====
|
||||
|
||||
struct NoopWaker;
|
||||
|
||||
fn noop_local_waker() -> LocalWaker {
|
||||
let w: NonNull<NoopWaker> = NonNull::dangling();
|
||||
unsafe { LocalWaker::new(w) }
|
||||
}
|
||||
|
||||
fn noop_waker() -> Waker {
|
||||
let w: NonNull<NoopWaker> = NonNull::dangling();
|
||||
unsafe { Waker::new(w) }
|
||||
}
|
||||
|
||||
unsafe impl UnsafeWake for NoopWaker {
|
||||
unsafe fn clone_raw(&self) -> Waker {
|
||||
noop_waker()
|
||||
}
|
||||
|
||||
unsafe fn drop_raw(&self) {
|
||||
}
|
||||
|
||||
unsafe fn wake(&self) {
|
||||
panic!("NoopWake cannot wake");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,68 @@
|
||||
|
||||
use futures::{Future, Async};
|
||||
|
||||
use std::marker::Unpin;
|
||||
use std::future::Future as StdFuture;
|
||||
use std::pin::Pin;
|
||||
use std::task::{LocalWaker, Poll as StdPoll};
|
||||
|
||||
/// Converts an 0.1 `Future` into an 0.3 `Future`.
|
||||
#[derive(Debug)]
|
||||
pub struct Compat<T>(T);
|
||||
|
||||
pub(crate) fn convert_poll<T, E>(poll: Result<Async<T>, E>) -> StdPoll<Result<T, E>> {
|
||||
use futures::Async::{Ready, NotReady};
|
||||
|
||||
match poll {
|
||||
Ok(Ready(val)) => StdPoll::Ready(Ok(val)),
|
||||
Ok(NotReady) => StdPoll::Pending,
|
||||
Err(err) => StdPoll::Ready(Err(err)),
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) fn convert_poll_stream<T, E>(
|
||||
poll: Result<Async<Option<T>>, E>) -> StdPoll<Option<Result<T, E>>>
|
||||
{
|
||||
use futures::Async::{Ready, NotReady};
|
||||
|
||||
match poll {
|
||||
Ok(Ready(Some(val))) => StdPoll::Ready(Some(Ok(val))),
|
||||
Ok(Ready(None)) => StdPoll::Ready(None),
|
||||
Ok(NotReady) => StdPoll::Pending,
|
||||
Err(err) => StdPoll::Ready(Some(Err(err))),
|
||||
}
|
||||
}
|
||||
|
||||
/// Convert a value into one that can be used with `await!`.
|
||||
pub trait IntoAwaitable {
|
||||
type Awaitable;
|
||||
|
||||
/// Convert `self` into a value that can be used with `await!`.
|
||||
fn into_awaitable(self) -> Self::Awaitable;
|
||||
}
|
||||
|
||||
impl<T: Future + Unpin> IntoAwaitable for T {
|
||||
type Awaitable = Compat<T>;
|
||||
|
||||
fn into_awaitable(self) -> Self::Awaitable {
|
||||
Compat(self)
|
||||
}
|
||||
}
|
||||
|
||||
impl<T> StdFuture for Compat<T>
|
||||
where T: Future + Unpin
|
||||
{
|
||||
type Output = Result<T::Item, T::Error>;
|
||||
|
||||
fn poll(mut self: Pin<&mut Self>, _lw: &LocalWaker) -> StdPoll<Self::Output> {
|
||||
use futures::Async::{Ready, NotReady};
|
||||
|
||||
// TODO: wire in cx
|
||||
|
||||
match self.0.poll() {
|
||||
Ok(Ready(val)) => StdPoll::Ready(Ok(val)),
|
||||
Ok(NotReady) => StdPoll::Pending,
|
||||
Err(e) => StdPoll::Ready(Err(e)),
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,11 +1,11 @@
|
||||
use tokio_io::AsyncWrite;
|
||||
|
||||
use futures_core::future::Future;
|
||||
use futures_core::task::{self, Poll};
|
||||
|
||||
use std::io;
|
||||
use std::future::Future;
|
||||
use std::marker::Unpin;
|
||||
use std::pin::PinMut;
|
||||
use std::pin::Pin;
|
||||
use std::task::{LocalWaker, Poll};
|
||||
|
||||
/// A future used to fully flush an I/O object.
|
||||
#[derive(Debug)]
|
||||
@@ -13,7 +13,7 @@ pub struct Flush<'a, T: ?Sized + 'a> {
|
||||
writer: &'a mut T,
|
||||
}
|
||||
|
||||
// PinMut is never projected to fields
|
||||
// Pin is never projected to fields
|
||||
impl<'a, T: ?Sized> Unpin for Flush<'a, T> {}
|
||||
|
||||
impl<'a, T: AsyncWrite + ?Sized> Flush<'a, T> {
|
||||
@@ -25,8 +25,8 @@ impl<'a, T: AsyncWrite + ?Sized> Flush<'a, T> {
|
||||
impl<'a, T: AsyncWrite + ?Sized> Future for Flush<'a, T> {
|
||||
type Output = io::Result<()>;
|
||||
|
||||
fn poll(mut self: PinMut<Self>, _cx: &mut task::Context) -> Poll<Self::Output> {
|
||||
use crate::async_await::compat::forward::convert_poll;
|
||||
fn poll(mut self: Pin<&mut Self>, _wx: &LocalWaker) -> Poll<Self::Output> {
|
||||
use crate::compat::forward::convert_poll;
|
||||
convert_poll(self.writer.poll_flush())
|
||||
}
|
||||
}
|
||||
@@ -1,11 +1,11 @@
|
||||
use tokio_io::AsyncRead;
|
||||
|
||||
use futures_core::future::Future;
|
||||
use futures_core::task::{self, Poll};
|
||||
use std::future::Future;
|
||||
use std::task::{self, Poll};
|
||||
|
||||
use std::io;
|
||||
use std::marker::Unpin;
|
||||
use std::pin::PinMut;
|
||||
use std::pin::Pin;
|
||||
|
||||
/// A future which can be used to read bytes.
|
||||
#[derive(Debug)]
|
||||
@@ -29,8 +29,8 @@ impl<'a, T: AsyncRead + ?Sized> Read<'a, T> {
|
||||
impl<'a, T: AsyncRead + ?Sized> Future for Read<'a, T> {
|
||||
type Output = io::Result<usize>;
|
||||
|
||||
fn poll(mut self: PinMut<Self>, _cx: &mut task::Context) -> Poll<Self::Output> {
|
||||
use crate::async_await::compat::forward::convert_poll;
|
||||
fn poll(mut self: Pin<&mut Self>, _lw: &task::LocalWaker) -> Poll<Self::Output> {
|
||||
use crate::compat::forward::convert_poll;
|
||||
|
||||
let this = &mut *self;
|
||||
convert_poll(this.reader.poll_read(this.buf))
|
||||
+5
-7
@@ -1,14 +1,12 @@
|
||||
use tokio_io::AsyncRead;
|
||||
|
||||
use futures_core::future::Future;
|
||||
use futures_core::task::{self, Poll};
|
||||
use futures_util::try_ready;
|
||||
use std::future::Future;
|
||||
use std::task::{self, Poll};
|
||||
|
||||
use std::io;
|
||||
use std::marker::Unpin;
|
||||
use std::mem;
|
||||
|
||||
use core::pin::PinMut;
|
||||
use std::pin::Pin;
|
||||
|
||||
/// A future which can be used to read exactly enough bytes to fill a buffer.
|
||||
#[derive(Debug)]
|
||||
@@ -36,8 +34,8 @@ fn eof() -> io::Error {
|
||||
impl<'a, T: AsyncRead + ?Sized> Future for ReadExact<'a, T> {
|
||||
type Output = io::Result<()>;
|
||||
|
||||
fn poll(mut self: PinMut<Self>, _cx: &mut task::Context) -> Poll<Self::Output> {
|
||||
use crate::async_await::compat::forward::convert_poll;
|
||||
fn poll(mut self: Pin<&mut Self>, _lw: &task::LocalWaker) -> Poll<Self::Output> {
|
||||
use crate::compat::forward::convert_poll;
|
||||
|
||||
let this = &mut *self;
|
||||
|
||||
@@ -1,11 +1,11 @@
|
||||
use tokio_io::AsyncWrite;
|
||||
|
||||
use futures_core::future::Future;
|
||||
use futures_core::task::{self, Poll};
|
||||
use std::future::Future;
|
||||
use std::task::{self, Poll};
|
||||
|
||||
use std::io;
|
||||
use std::marker::Unpin;
|
||||
use std::pin::PinMut;
|
||||
use std::pin::Pin;
|
||||
|
||||
/// A future used to write data.
|
||||
#[derive(Debug)]
|
||||
@@ -29,8 +29,8 @@ impl<'a, T: AsyncWrite + ?Sized> Write<'a, T> {
|
||||
impl<'a, T: AsyncWrite + ?Sized> Future for Write<'a, T> {
|
||||
type Output = io::Result<usize>;
|
||||
|
||||
fn poll(mut self: PinMut<Self>, _cx: &mut task::Context) -> Poll<io::Result<usize>> {
|
||||
use crate::async_await::compat::forward::convert_poll;
|
||||
fn poll(mut self: Pin<&mut Self>, _lw: &task::LocalWaker) -> Poll<io::Result<usize>> {
|
||||
use crate::compat::forward::convert_poll;
|
||||
|
||||
let this = &mut *self;
|
||||
convert_poll(this.writer.poll_write(this.buf))
|
||||
+5
-7
@@ -1,14 +1,12 @@
|
||||
use tokio_io::AsyncWrite;
|
||||
|
||||
use futures_core::future::Future;
|
||||
use futures_core::task::{self, Poll};
|
||||
use futures_util::try_ready;
|
||||
use std::future::Future;
|
||||
use std::task::{self, Poll};
|
||||
|
||||
use std::io;
|
||||
use std::marker::Unpin;
|
||||
use std::mem;
|
||||
|
||||
use core::pin::PinMut;
|
||||
use std::pin::Pin;
|
||||
|
||||
/// A future used to write the entire contents of a buffer.
|
||||
#[derive(Debug)]
|
||||
@@ -36,8 +34,8 @@ fn zero_write() -> io::Error {
|
||||
impl<'a, T: AsyncWrite + ?Sized> Future for WriteAll<'a, T> {
|
||||
type Output = io::Result<()>;
|
||||
|
||||
fn poll(mut self: PinMut<Self>, _cx: &mut task::Context) -> Poll<io::Result<()>> {
|
||||
use crate::async_await::compat::forward::convert_poll;
|
||||
fn poll(mut self: Pin<&mut Self>, _lw: &task::LocalWaker) -> Poll<io::Result<()>> {
|
||||
use crate::compat::forward::convert_poll;
|
||||
|
||||
let this = &mut *self;
|
||||
|
||||
@@ -1,4 +1,12 @@
|
||||
#![feature(futures_api, await_macro, pin, arbitrary_self_types)]
|
||||
#![cfg(feature = "async-await-preview")]
|
||||
#![feature(
|
||||
rust_2018_preview,
|
||||
arbitrary_self_types,
|
||||
async_await,
|
||||
await_macro,
|
||||
futures_api,
|
||||
pin,
|
||||
)]
|
||||
|
||||
#![doc(html_root_url = "https://docs.rs/tokio-async-await/0.1.3")]
|
||||
#![deny(missing_docs, missing_debug_implementations)]
|
||||
@@ -7,39 +15,31 @@
|
||||
//! A preview of Tokio w/ `async` / `await` support.
|
||||
|
||||
extern crate futures;
|
||||
extern crate futures_core;
|
||||
extern crate futures_util;
|
||||
extern crate tokio_io;
|
||||
|
||||
// Re-export all of Tokio
|
||||
pub use tokio_main::{
|
||||
// Modules
|
||||
clock,
|
||||
codec,
|
||||
executor,
|
||||
fs,
|
||||
io,
|
||||
net,
|
||||
reactor,
|
||||
runtime,
|
||||
timer,
|
||||
util,
|
||||
|
||||
// Functions
|
||||
run,
|
||||
spawn,
|
||||
};
|
||||
|
||||
pub mod sync {
|
||||
//! Asynchronous aware synchronization
|
||||
|
||||
pub use tokio_channel::{
|
||||
mpsc,
|
||||
oneshot,
|
||||
};
|
||||
/// Extracts the successful type of a `Poll<Result<T, E>>`.
|
||||
///
|
||||
/// This macro bakes in propagation of `Pending` and `Err` signals by returning early.
|
||||
macro_rules! try_ready {
|
||||
($x:expr) => {
|
||||
match $x {
|
||||
std::task::Poll::Ready(Ok(x)) => x,
|
||||
std::task::Poll::Ready(Err(e)) =>
|
||||
return std::task::Poll::Ready(Err(e.into())),
|
||||
std::task::Poll::Pending =>
|
||||
return std::task::Poll::Pending,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub mod async_await;
|
||||
#[macro_use]
|
||||
mod await;
|
||||
pub mod compat;
|
||||
pub mod io;
|
||||
pub mod sink;
|
||||
pub mod stream;
|
||||
|
||||
/*
|
||||
pub mod prelude {
|
||||
//! A "prelude" for users of the `tokio` crate.
|
||||
//!
|
||||
@@ -69,33 +69,47 @@ pub mod prelude {
|
||||
},
|
||||
};
|
||||
}
|
||||
*/
|
||||
|
||||
use futures_core::{
|
||||
Future as Future03,
|
||||
};
|
||||
|
||||
// Rename the `await` macro in `std`
|
||||
// Rename the `await` macro in `std`. This is used by the redefined
|
||||
// `await` macro in this crate.
|
||||
#[doc(hidden)]
|
||||
pub use std::await as std_await;
|
||||
|
||||
/*
|
||||
use std::future::{Future as StdFuture};
|
||||
|
||||
fn run<T: futures::Future<Item = (), Error = ()>>(t: T) {
|
||||
drop(t);
|
||||
}
|
||||
|
||||
async fn map_ok<T: StdFuture>(future: T) -> Result<(), ()> {
|
||||
let _ = await!(future);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Like `tokio::run`, but takes an `async` block
|
||||
pub fn run_async<F>(future: F)
|
||||
where F: Future03<Output = ()> + Send + 'static,
|
||||
where F: StdFuture<Output = ()> + Send + 'static,
|
||||
{
|
||||
use futures_util::future::FutureExt;
|
||||
use crate::async_await::compat::backward;
|
||||
use async_await::compat::backward;
|
||||
let future = backward::Compat::new(map_ok(future));
|
||||
|
||||
let future = future.map(|_| Ok(()));
|
||||
run(backward::Compat::new(future))
|
||||
run(future);
|
||||
unimplemented!();
|
||||
}
|
||||
*/
|
||||
|
||||
/*
|
||||
/// Like `tokio::spawn`, but takes an `async` block
|
||||
pub fn spawn_async<F>(future: F)
|
||||
where F: Future03<Output = ()> + Send + 'static,
|
||||
where F: StdFuture<Output = ()> + Send + 'static,
|
||||
{
|
||||
use futures_util::future::FutureExt;
|
||||
use crate::async_await::compat::backward;
|
||||
|
||||
let future = future.map(|_| Ok(()));
|
||||
spawn(backward::Compat::new(future));
|
||||
spawn(backward::Compat::new(async || {
|
||||
let _ = await!(future);
|
||||
Ok(())
|
||||
}));
|
||||
}
|
||||
*/
|
||||
|
||||
+8
-13
@@ -1,10 +1,10 @@
|
||||
use futures::Sink;
|
||||
|
||||
use futures_core::future::Future;
|
||||
use futures_core::task::{self, Poll};
|
||||
use std::future::Future;
|
||||
use std::task::{self, Poll};
|
||||
|
||||
use std::marker::Unpin;
|
||||
use std::pin::PinMut;
|
||||
use std::pin::Pin;
|
||||
|
||||
/// Future for the `SinkExt::send_async` combinator, which sends a value to a
|
||||
/// sink and then waits until the sink has fully flushed.
|
||||
@@ -28,17 +28,12 @@ impl<'a, T: Sink + Unpin + ?Sized> Send<'a, T> {
|
||||
impl<T: Sink + Unpin + ?Sized> Future for Send<'_, T> {
|
||||
type Output = Result<(), T::SinkError>;
|
||||
|
||||
fn poll(mut self: PinMut<Self>, _cx: &mut task::Context) -> Poll<Self::Output> {
|
||||
use crate::async_await::compat::forward::convert_poll;
|
||||
fn poll(mut self: Pin<&mut Self>, _lw: &task::LocalWaker) -> Poll<Self::Output> {
|
||||
use crate::compat::forward::convert_poll;
|
||||
use futures::AsyncSink::{Ready, NotReady};
|
||||
use futures_util::try_ready;
|
||||
|
||||
// use crate::compat::forward::convert_poll;
|
||||
|
||||
let this = &mut *self;
|
||||
|
||||
if let Some(item) = this.item.take() {
|
||||
match this.sink.start_send(item) {
|
||||
if let Some(item) = self.item.take() {
|
||||
match self.sink.start_send(item) {
|
||||
Ok(Ready) => {}
|
||||
Ok(NotReady(val)) => {
|
||||
self.item = Some(val);
|
||||
@@ -52,7 +47,7 @@ impl<T: Sink + Unpin + ?Sized> Future for Send<'_, T> {
|
||||
|
||||
// we're done sending the item, but want to block on flushing the
|
||||
// sink
|
||||
try_ready!(convert_poll(this.sink.poll_complete()));
|
||||
try_ready!(convert_poll(self.sink.poll_complete()));
|
||||
|
||||
Poll::Ready(Ok(()))
|
||||
}
|
||||
+6
-7
@@ -1,9 +1,9 @@
|
||||
use futures::Stream;
|
||||
use futures_core::future::Future;
|
||||
use futures_core::task::{self, Poll};
|
||||
|
||||
use std::future::Future;
|
||||
use std::marker::Unpin;
|
||||
use std::pin::PinMut;
|
||||
use std::pin::Pin;
|
||||
use std::task::{LocalWaker, Poll};
|
||||
|
||||
/// A future of the next element of a stream.
|
||||
#[derive(Debug)]
|
||||
@@ -22,10 +22,9 @@ impl<'a, T: Stream + Unpin> Next<'a, T> {
|
||||
impl<'a, T: Stream + Unpin> Future for Next<'a, T> {
|
||||
type Output = Option<Result<T::Item, T::Error>>;
|
||||
|
||||
fn poll(self: PinMut<Self>, _cx: &mut task::Context) -> Poll<Self::Output> {
|
||||
use crate::async_await::compat::forward::convert_poll_stream;
|
||||
fn poll(mut self: Pin<&mut Self>, _lw: &LocalWaker) -> Poll<Self::Output> {
|
||||
use crate::compat::forward::convert_poll_stream;
|
||||
|
||||
convert_poll_stream(
|
||||
PinMut::get_mut(self).stream.poll())
|
||||
convert_poll_stream(self.stream.poll())
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user