New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Allow configuring a timeout for Endpoint
selection in EndpointGroup
#4246
Changes from 11 commits
9e50dfd
44b78e7
262935e
5837649
afb526d
9030415
a34ce58
c399ece
af45d0b
09afaef
da308f5
2a6d00b
cc09bf6
4b1df0d
409a3b2
1ffbeb8
eba5255
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,46 @@ | ||
/* | ||
* Copyright 2022 LINE Corporation | ||
* | ||
* LINE Corporation licenses this file to you under the Apache License, | ||
* version 2.0 (the "License"); you may not use this file except in compliance | ||
* with the License. You may obtain a copy of the License at: | ||
* | ||
* https://www.apache.org/licenses/LICENSE-2.0 | ||
* | ||
* Unless required by applicable law or agreed to in writing, software | ||
* distributed under the License is distributed on an "AS IS" BASIS, WITHOUT | ||
* WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the | ||
* License for the specific language governing permissions and limitations | ||
* under the License. | ||
*/ | ||
|
||
package com.linecorp.armeria.client.consul; | ||
|
||
import static org.assertj.core.api.Assertions.assertThat; | ||
|
||
import java.net.URI; | ||
|
||
import org.junit.jupiter.api.Test; | ||
|
||
import com.linecorp.armeria.common.Flags; | ||
|
||
class ConsulEndpointGroupBuilderTest { | ||
|
||
@Test | ||
void selectionTimeout_default() { | ||
try (ConsulEndpointGroup group = ConsulEndpointGroup.of(URI.create("http://127.0.0.1/node"), | ||
"my-service")) { | ||
assertThat(group.selectionTimeoutMillis()).isEqualTo(Flags.defaultResponseTimeoutMillis()); | ||
} | ||
} | ||
|
||
@Test | ||
void selectionTimeout_custom() { | ||
try (ConsulEndpointGroup group = | ||
ConsulEndpointGroup.builder(URI.create("http://127.0.0.1/node"), "my-service") | ||
.selectionTimeoutMillis(4000) | ||
.build()) { | ||
assertThat(group.selectionTimeoutMillis()).isEqualTo(4000); | ||
} | ||
} | ||
} |
Original file line number | Diff line number | Diff line change |
---|---|---|
|
@@ -25,14 +25,16 @@ | |
import java.util.concurrent.TimeUnit; | ||
import java.util.function.Consumer; | ||
|
||
import com.google.common.annotations.VisibleForTesting; | ||
|
||
import com.linecorp.armeria.client.ClientRequestContext; | ||
import com.linecorp.armeria.client.Endpoint; | ||
import com.linecorp.armeria.common.annotation.Nullable; | ||
import com.linecorp.armeria.common.util.UnmodifiableFuture; | ||
|
||
/** | ||
* A skeletal {@link EndpointSelector} implementation. This abstract class implements the | ||
* {@link #select(ClientRequestContext, ScheduledExecutorService, long)} method by listening to | ||
* {@link #select(ClientRequestContext, ScheduledExecutorService)} method by listening to | ||
* the change events emitted by {@link EndpointGroup} specified at construction time. | ||
*/ | ||
public abstract class AbstractEndpointSelector implements EndpointSelector { | ||
|
@@ -53,10 +55,17 @@ protected final EndpointGroup group() { | |
return endpointGroup; | ||
} | ||
|
||
@Deprecated | ||
@Override | ||
public final CompletableFuture<Endpoint> select(ClientRequestContext ctx, | ||
ScheduledExecutorService executor, | ||
long timeoutMillis) { | ||
return select(ctx, executor); | ||
} | ||
|
||
@Override | ||
public final CompletableFuture<Endpoint> select(ClientRequestContext ctx, | ||
ScheduledExecutorService executor) { | ||
Endpoint endpoint = selectNow(ctx); | ||
if (endpoint != null) { | ||
return UnmodifiableFuture.completedFuture(endpoint); | ||
|
@@ -73,25 +82,41 @@ public final CompletableFuture<Endpoint> select(ClientRequestContext ctx, | |
return UnmodifiableFuture.completedFuture(endpoint); | ||
} | ||
|
||
final long selectionTimeoutMillis = endpointGroup.selectionTimeoutMillis(); | ||
if (selectionTimeoutMillis == 0) { | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. nit: Would it be better to use There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. I don't think we need to differentiate the timeout of |
||
// A static EndpointGroup. | ||
return UnmodifiableFuture.completedFuture(null); | ||
} | ||
long responseTimeoutMillis = ctx.responseTimeoutMillis(); | ||
if (responseTimeoutMillis == 0) { | ||
responseTimeoutMillis = Long.MAX_VALUE; | ||
} | ||
|
||
final long timeoutMillis = Math.min(selectionTimeoutMillis, responseTimeoutMillis); | ||
|
||
// Schedule the timeout task. | ||
final ScheduledFuture<?> timeoutFuture = | ||
executor.schedule(() -> listeningFuture.complete(null), | ||
timeoutMillis, TimeUnit.MILLISECONDS); | ||
listeningFuture.timeoutFuture = timeoutFuture; | ||
|
||
// Cancel the timeout task if listeningFuture is done already. | ||
// This guards against the following race condition: | ||
// 1) (Current thread) Timeout task is scheduled. | ||
// 2) ( Other thread ) listeningFuture is completed, but the timeout task is not cancelled | ||
// 3) (Current thread) timeoutFuture is assigned to listeningFuture.timeoutFuture, but it's too late. | ||
if (listeningFuture.isDone()) { | ||
timeoutFuture.cancel(false); | ||
if (timeoutMillis < Long.MAX_VALUE) { | ||
final ScheduledFuture<?> timeoutFuture = executor.schedule(() -> { | ||
listeningFuture.complete(null); | ||
}, timeoutMillis, TimeUnit.MILLISECONDS); | ||
listeningFuture.timeoutFuture = timeoutFuture; | ||
|
||
// Cancel the timeout task if listeningFuture is done already. | ||
// This guards against the following race condition: | ||
// 1) (Current thread) Timeout task is scheduled. | ||
// 2) ( Other thread ) listeningFuture is completed, but the timeout task is not cancelled | ||
// 3) (Current thread) timeoutFuture is assigned to listeningFuture.timeoutFuture, but it's too | ||
// late. | ||
if (listeningFuture.isDone()) { | ||
timeoutFuture.cancel(false); | ||
} | ||
} | ||
|
||
return listeningFuture; | ||
} | ||
|
||
private class ListeningFuture extends CompletableFuture<Endpoint> implements Consumer<List<Endpoint>> { | ||
@VisibleForTesting | ||
class ListeningFuture extends CompletableFuture<Endpoint> implements Consumer<List<Endpoint>> { | ||
ikhoon marked this conversation as resolved.
Show resolved
Hide resolved
|
||
private final ClientRequestContext ctx; | ||
private final Executor executor; | ||
@Nullable | ||
|
@@ -150,5 +175,11 @@ private void cleanup() { | |
timeoutFuture.cancel(false); | ||
} | ||
} | ||
|
||
@Nullable | ||
@VisibleForTesting | ||
ScheduledFuture<?> timeoutFuture() { | ||
return timeoutFuture; | ||
} | ||
} | ||
} |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
I think the default for
DnsEndpointGroupBuilder
will beconnectTimeout
instead ofresponseTimeout
. I wanted to check if this is your intention.If the intention is to set
responseTimeout
as the default for pre-definedEndpointGroup
s inarmeria
, what do you think of just removing the empty constructor?We can catch these type of mistakes at compile-time this way.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Sorry for the late response.
Currently, the default select timeout for
DnsEndpointGroup
is the default connection timeout.armeria/core/src/test/java/com/linecorp/armeria/client/endpoint/SelectionTimeoutTest.java
Line 118 in af45d0b
DNS servers usually return a response quickly because they cache DNS records. Therefore, 3.2 seconds may be reasonable.
By the way, I checked the default timeout for a DNS query. It is 5 seconds that can be a more sensible default timeout for `DnsEndpointGroup/
As this class is a public API, that could cause a breaking change.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
AbstractDynamicEndpointGroupBuilder is part of unstable APIs. Will remove the default constructor.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
It sounds reasonable to use the same timeout for
DnsEndpointGroup
and our DNS resolver.