fix(transport): return Poll::ready until error is consumed (#536)
* fix(transport): return Poll::ready until error is consumed When a lazy connection fails to connect it first returns Poll::ready from the reconnect service, yet the subsequent call returns Poll::pending making tower_balance loop forever. Instead, on error we return Ready until the error is consumed in the call method. * chore: revert version change * refactor: Into<Error> bounds for error intead of debug * Remove fmt::Debug bound for reconnect Co-authored-by: Helge Hoff <[email protected]>
This commit is contained in:
co-authored by
Helge Hoff
parent
31936e0513
commit
dafea9adee
@@ -13,11 +13,12 @@ use tracing::trace;
|
|||||||
pub(crate) struct Reconnect<M, Target>
|
pub(crate) struct Reconnect<M, Target>
|
||||||
where
|
where
|
||||||
M: Service<Target>,
|
M: Service<Target>,
|
||||||
|
M::Error: Into<Error>,
|
||||||
{
|
{
|
||||||
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>,
|
error: Option<crate::Error>,
|
||||||
has_been_connected: bool,
|
has_been_connected: bool,
|
||||||
is_lazy: bool,
|
is_lazy: bool,
|
||||||
}
|
}
|
||||||
@@ -32,6 +33,7 @@ enum State<F, S> {
|
|||||||
impl<M, Target> Reconnect<M, Target>
|
impl<M, Target> Reconnect<M, Target>
|
||||||
where
|
where
|
||||||
M: Service<Target>,
|
M: Service<Target>,
|
||||||
|
M::Error: Into<Error>,
|
||||||
{
|
{
|
||||||
pub(crate) fn new(mk_service: M, target: Target, is_lazy: bool) -> Self {
|
pub(crate) fn new(mk_service: M, target: Target, is_lazy: bool) -> Self {
|
||||||
Reconnect {
|
Reconnect {
|
||||||
@@ -52,14 +54,19 @@ where
|
|||||||
M::Future: Unpin,
|
M::Future: Unpin,
|
||||||
Error: From<M::Error> + From<S::Error>,
|
Error: From<M::Error> + From<S::Error>,
|
||||||
Target: Clone,
|
Target: Clone,
|
||||||
|
<M as tower_service::Service<Target>>::Error: Into<crate::Error>,
|
||||||
{
|
{
|
||||||
type Response = S::Response;
|
type Response = S::Response;
|
||||||
type Error = Error;
|
type Error = Error;
|
||||||
type Future = ResponseFuture<S::Future, M::Error>;
|
type Future = ResponseFuture<S::Future>;
|
||||||
|
|
||||||
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 mut state;
|
let mut state;
|
||||||
|
|
||||||
|
if self.error.is_some() {
|
||||||
|
return Poll::Ready(Ok(()));
|
||||||
|
}
|
||||||
|
|
||||||
loop {
|
loop {
|
||||||
match self.state {
|
match self.state {
|
||||||
State::Idle => {
|
State::Idle => {
|
||||||
@@ -94,7 +101,9 @@ where
|
|||||||
if !(self.has_been_connected || self.is_lazy) {
|
if !(self.has_been_connected || self.is_lazy) {
|
||||||
return Poll::Ready(Err(e.into()));
|
return Poll::Ready(Err(e.into()));
|
||||||
} else {
|
} else {
|
||||||
self.error = Some(e);
|
let error = e.into();
|
||||||
|
tracing::error!("reconnect::poll_ready: {:?}", error);
|
||||||
|
self.error = Some(error);
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -130,8 +139,10 @@ where
|
|||||||
}
|
}
|
||||||
|
|
||||||
fn call(&mut self, request: Request) -> Self::Future {
|
fn call(&mut self, request: Request) -> Self::Future {
|
||||||
|
tracing::trace!("Reconnect::call");
|
||||||
if let Some(error) = self.error.take() {
|
if let Some(error) = self.error.take() {
|
||||||
return ResponseFuture::error(error);
|
tracing::error!("error: {:?}", error);
|
||||||
|
return ResponseFuture::error(error.into());
|
||||||
}
|
}
|
||||||
|
|
||||||
let service = match self.state {
|
let service = match self.state {
|
||||||
@@ -150,6 +161,7 @@ where
|
|||||||
M::Future: fmt::Debug,
|
M::Future: fmt::Debug,
|
||||||
M::Response: fmt::Debug,
|
M::Response: fmt::Debug,
|
||||||
Target: fmt::Debug,
|
Target: fmt::Debug,
|
||||||
|
<M as tower_service::Service<Target>>::Error: Into<Error>,
|
||||||
{
|
{
|
||||||
fn fmt(&self, fmt: &mut fmt::Formatter<'_>) -> fmt::Result {
|
fn fmt(&self, fmt: &mut fmt::Formatter<'_>) -> fmt::Result {
|
||||||
fmt.debug_struct("Reconnect")
|
fmt.debug_struct("Reconnect")
|
||||||
@@ -163,37 +175,36 @@ where
|
|||||||
/// Future that resolves to the response or failure to connect.
|
/// Future that resolves to the response or failure to connect.
|
||||||
#[pin_project]
|
#[pin_project]
|
||||||
#[derive(Debug)]
|
#[derive(Debug)]
|
||||||
pub(crate) struct ResponseFuture<F, E> {
|
pub(crate) struct ResponseFuture<F> {
|
||||||
#[pin]
|
#[pin]
|
||||||
inner: Inner<F, E>,
|
inner: Inner<F>,
|
||||||
}
|
}
|
||||||
|
|
||||||
#[pin_project(project = InnerProj)]
|
#[pin_project(project = InnerProj)]
|
||||||
#[derive(Debug)]
|
#[derive(Debug)]
|
||||||
enum Inner<F, E> {
|
enum Inner<F> {
|
||||||
Future(#[pin] F),
|
Future(#[pin] F),
|
||||||
Error(Option<E>),
|
Error(Option<crate::Error>),
|
||||||
}
|
}
|
||||||
|
|
||||||
impl<F, E> ResponseFuture<F, E> {
|
impl<F> ResponseFuture<F> {
|
||||||
pub(crate) fn new(inner: F) -> Self {
|
pub(crate) fn new(inner: F) -> Self {
|
||||||
ResponseFuture {
|
ResponseFuture {
|
||||||
inner: Inner::Future(inner),
|
inner: Inner::Future(inner),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
pub(crate) fn error(error: E) -> Self {
|
pub(crate) fn error(error: crate::Error) -> Self {
|
||||||
ResponseFuture {
|
ResponseFuture {
|
||||||
inner: Inner::Error(Some(error)),
|
inner: Inner::Error(Some(error)),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
impl<F, T, E, ME> Future for ResponseFuture<F, ME>
|
impl<F, T, E> Future for ResponseFuture<F>
|
||||||
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>;
|
||||||
|
|
||||||
@@ -203,7 +214,7 @@ where
|
|||||||
match me.inner.project() {
|
match me.inner.project() {
|
||||||
InnerProj::Future(fut) => fut.poll(cx).map_err(Into::into),
|
InnerProj::Future(fut) => fut.poll(cx).map_err(Into::into),
|
||||||
InnerProj::Error(e) => {
|
InnerProj::Error(e) => {
|
||||||
let e = e.take().expect("Polled after ready.").into();
|
let e = e.take().expect("Polled after ready.");
|
||||||
Poll::Ready(Err(e))
|
Poll::Ready(Err(e))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user