fix(transport): Fix lazily reconnecting (#187)

Closes #167
This commit is contained in:
Lucio Franco
2019-12-13 20:24:45 -05:00
committed by GitHub
parent 97d5363e2b
commit 0505dff65a
+43 -11
View File
@@ -1,5 +1,5 @@
use crate::Error; use crate::Error;
use pin_project::pin_project; use pin_project::{pin_project, project};
use std::fmt; use std::fmt;
use std::{ use std::{
future::Future, future::Future,
@@ -17,6 +17,7 @@ where
mk_service: M, mk_service: M,
state: State<M::Future, M::Response>, state: State<M::Future, M::Response>,
target: Target, target: Target,
error: Option<M::Error>,
} }
#[derive(Debug)] #[derive(Debug)]
@@ -41,6 +42,7 @@ where
mk_service, mk_service,
state: State::Connected(initial_connection), state: State::Connected(initial_connection),
target, target,
error: None,
} }
} }
} }
@@ -55,10 +57,9 @@ where
{ {
type Response = S::Response; type Response = S::Response;
type Error = Error; type Error = Error;
type Future = ResponseFuture<S::Future>; type Future = ResponseFuture<S::Future, M::Error>;
fn poll_ready(&mut self, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> { fn poll_ready(&mut self, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
let ret;
let mut state; let mut state;
loop { loop {
@@ -90,7 +91,7 @@ where
Poll::Ready(Err(e)) => { Poll::Ready(Err(e)) => {
trace!("poll_ready; error"); trace!("poll_ready; error");
state = State::Idle; state = State::Idle;
ret = Err(e.into()); self.error = Some(e.into());
break; break;
} }
} }
@@ -118,10 +119,14 @@ where
} }
self.state = state; self.state = state;
Poll::Ready(ret) Poll::Ready(Ok(()))
} }
fn call(&mut self, request: Request) -> Self::Future { fn call(&mut self, request: Request) -> Self::Future {
if let Some(error) = self.error.take() {
return ResponseFuture::error(error);
}
let service = match self.state { let service = match self.state {
State::Connected(ref mut service) => service, State::Connected(ref mut service) => service,
_ => panic!("service not ready; poll_ready must be called first"), _ => panic!("service not ready; poll_ready must be called first"),
@@ -148,27 +153,54 @@ where
} }
} }
/// Future that resolves to the response or failure to connect.
#[pin_project] #[pin_project]
#[derive(Debug)] #[derive(Debug)]
pub(crate) struct ResponseFuture<F> { pub(crate) struct ResponseFuture<F, E> {
#[pin] #[pin]
inner: F, inner: Inner<F, E>,
} }
impl<F> ResponseFuture<F> { #[pin_project]
#[derive(Debug)]
enum Inner<F, E> {
Future(#[pin] F),
Error(Option<E>),
}
impl<F, E> ResponseFuture<F, E> {
pub(crate) fn new(inner: F) -> Self { pub(crate) fn new(inner: F) -> Self {
ResponseFuture { inner } ResponseFuture {
inner: Inner::Future(inner),
}
}
pub(crate) fn error(error: E) -> Self {
ResponseFuture {
inner: Inner::Error(Some(error)),
}
} }
} }
impl<F, T, E> Future for ResponseFuture<F> impl<F, T, E, ME> Future for ResponseFuture<F, ME>
where where
F: Future<Output = Result<T, E>>, F: Future<Output = Result<T, E>>,
E: Into<Error>, E: Into<Error>,
ME: Into<Error>,
{ {
type Output = Result<T, Error>; type Output = Result<T, Error>;
#[project]
fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> { fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
self.project().inner.poll(cx).map_err(Into::into) //self.project().inner.poll(cx).map_err(Into::into)
let me = self.project();
#[project]
match me.inner.project() {
Inner::Future(fut) => fut.poll(cx).map_err(Into::into),
Inner::Error(e) => {
let e = e.take().expect("Polled after ready.").into();
Poll::Ready(Err(e))
}
}
} }
} }