use futures::{future, prelude::*};
pub fn send<S: Sink>(s: S, x: S::SinkItem) -> impl Future<Item = S, Error = S::SinkError> {
SendFuture { sink: Some(s), item: Some(x) }
}
pub fn flush<S: Sink>(s: S) -> impl Future<Item = S, Error = S::SinkError> {
s.flush()
}
pub fn close<S: Sink>(mut s: S) -> impl Future<Item = (), Error = S::SinkError> {
future::poll_fn(move || s.close())
}
struct SendFuture<S: Sink> {
sink: Option<S>,
item: Option<S::SinkItem>
}
impl<S: Sink> Future for SendFuture<S> {
type Item = S;
type Error = S::SinkError;
fn poll(&mut self) -> Poll<S, S::SinkError> {
let item = self.item.take().expect("SendFuture has not completed => item exists");
let mut sink = self.sink.take().expect("SendFuture has not completed => sink exists");
match sink.start_send(item)? {
AsyncSink::Ready => Ok(Async::Ready(sink)),
AsyncSink::NotReady(item) => {
self.item = Some(item);
self.sink = Some(sink);
Ok(Async::NotReady)
}
}
}
}