Skip to content

Commit e0d43b3

Browse files
authored
xDS retry now correctly re-selects endpoints on each retry attempt. (#6798)
1 parent 423d994 commit e0d43b3

18 files changed

Lines changed: 422 additions & 185 deletions

File tree

it/xds-client/src/test/java/com/linecorp/armeria/xds/it/PreprocessorErrorTest.java

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,7 @@
3232
import com.linecorp.armeria.client.Clients;
3333
import com.linecorp.armeria.client.UnprocessedRequestException;
3434
import com.linecorp.armeria.client.WebClient;
35+
import com.linecorp.armeria.client.endpoint.EmptyEndpointGroupException;
3536
import com.linecorp.armeria.common.TimeoutException;
3637
import com.linecorp.armeria.internal.client.DefaultClientRequestContext;
3738
import com.linecorp.armeria.xds.XdsBootstrap;
@@ -133,8 +134,8 @@ static Stream<Arguments> testCases() {
133134
endpoints:
134135
- lb_endpoints:
135136
""",
136-
TimeoutException.class,
137-
"Failed to select an endpoint"
137+
EmptyEndpointGroupException.class,
138+
"Unable to select endpoints from"
138139
)
139140
);
140141
}

it/xds-client/src/test/java/com/linecorp/armeria/xds/it/RetryTest.java

Lines changed: 92 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,16 +22,20 @@
2222

2323
import java.util.ArrayDeque;
2424
import java.util.List;
25+
import java.util.Set;
2526
import java.util.concurrent.TimeUnit;
27+
import java.util.stream.Collectors;
2628
import java.util.stream.Stream;
2729

30+
import org.junit.jupiter.api.Test;
2831
import org.junit.jupiter.params.ParameterizedTest;
2932
import org.junit.jupiter.params.provider.Arguments;
3033
import org.junit.jupiter.params.provider.MethodSource;
3134

3235
import com.linecorp.armeria.client.ClientRequestContext;
3336
import com.linecorp.armeria.client.ClientRequestContextCaptor;
3437
import com.linecorp.armeria.client.Clients;
38+
import com.linecorp.armeria.client.Endpoint;
3539
import com.linecorp.armeria.client.RefusedStreamException;
3640
import com.linecorp.armeria.client.UnprocessedRequestException;
3741
import com.linecorp.armeria.client.WebClient;
@@ -774,4 +778,92 @@ void retriableRequestHeaders(String retryOptions, ResponseHeaders responseHeader
774778
assertThat(ctx.log().children()).hasSize(expectedRetries + 1);
775779
}
776780
}
781+
782+
//language=YAML
783+
private static final String multiEndpointBootstrap =
784+
"""
785+
static_resources:
786+
listeners:
787+
- name: my-listener
788+
api_listener:
789+
api_listener:
790+
"@type": type.googleapis.com/envoy.extensions.filters.network.\
791+
http_connection_manager.v3.HttpConnectionManager
792+
stat_prefix: http
793+
route_config:
794+
name: local_route
795+
virtual_hosts:
796+
- name: local_service1
797+
domains: [ "*" ]
798+
routes:
799+
- match:
800+
prefix: /
801+
route:
802+
cluster: my-cluster
803+
retry_policy:
804+
retry_on: "5xx"
805+
num_retries: 3
806+
http_filters:
807+
- name: envoy.filters.http.router
808+
typed_config:
809+
"@type": type.googleapis.com/envoy.extensions.filters.http.router.v3.Router
810+
clusters:
811+
- name: my-cluster
812+
type: STATIC
813+
load_assignment:
814+
cluster_name: my-cluster
815+
endpoints:
816+
- lb_endpoints:
817+
- endpoint:
818+
address:
819+
socket_address:
820+
address: 127.0.0.1
821+
port_value: 8080
822+
- endpoint:
823+
address:
824+
socket_address:
825+
address: 127.0.0.1
826+
port_value: 8081
827+
- endpoint:
828+
address:
829+
socket_address:
830+
address: 127.0.0.1
831+
port_value: 8082
832+
- endpoint:
833+
address:
834+
socket_address:
835+
address: 127.0.0.1
836+
port_value: 8083
837+
""";
838+
839+
@Test
840+
void retrySelectsDifferentEndpoints() {
841+
try (XdsBootstrap xdsBootstrap = XdsBootstrap.of(XdsResourceReader.fromYaml(multiEndpointBootstrap));
842+
XdsHttpPreprocessor preprocessor = XdsHttpPreprocessor.ofListener("my-listener", xdsBootstrap)) {
843+
final ArrayDeque<Endpoint> selectedEndpoints = new ArrayDeque<>();
844+
final ClientRequestContext ctx;
845+
try (ClientRequestContextCaptor captor = Clients.newContextCaptor()) {
846+
final AggregatedHttpResponse res =
847+
WebClient.builder(preprocessor)
848+
.decorator((delegate, ctx0, req) -> {
849+
selectedEndpoints.add(ctx0.endpoint());
850+
return HttpResponse.of(ResponseHeaders.of(503));
851+
})
852+
.build()
853+
.blocking()
854+
.execute(HttpRequest.of(HttpMethod.GET, "/"));
855+
assertThat(res.status().code()).isEqualTo(503);
856+
ctx = captor.get();
857+
}
858+
// 1 original + 3 retries = 4 total attempts
859+
assertThat(ctx.log().children()).hasSize(4);
860+
assertThat(selectedEndpoints).hasSize(4);
861+
862+
// Verify that not all attempts used the same endpoint
863+
final Set<Integer> uniquePorts = selectedEndpoints.stream()
864+
.map(Endpoint::port)
865+
.collect(Collectors.toSet());
866+
assertThat(uniquePorts).hasSizeGreaterThan(1);
867+
}
868+
}
777869
}
Lines changed: 148 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,148 @@
1+
/*
2+
* Copyright 2025 LY Corporation
3+
*
4+
* LY Corporation licenses this file to you under the Apache License,
5+
* version 2.0 (the "License"); you may not use this file except in compliance
6+
* with the License. You may obtain a copy of the License at:
7+
*
8+
* https://www.apache.org/licenses/LICENSE-2.0
9+
*
10+
* Unless required by applicable law or agreed to in writing, software
11+
* distributed under the License is distributed on an "AS IS" BASIS, WITHOUT
12+
* WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the
13+
* License for the specific language governing permissions and limitations
14+
* under the License.
15+
*/
16+
17+
package com.linecorp.armeria.xds;
18+
19+
import java.util.Set;
20+
21+
import com.linecorp.armeria.client.ClientDecoration;
22+
import com.linecorp.armeria.client.ClientRequestContext;
23+
import com.linecorp.armeria.client.ClientTlsSpec;
24+
import com.linecorp.armeria.client.Endpoint;
25+
import com.linecorp.armeria.client.HttpClient;
26+
import com.linecorp.armeria.client.HttpPreprocessor;
27+
import com.linecorp.armeria.client.PreClient;
28+
import com.linecorp.armeria.client.PreClientRequestContext;
29+
import com.linecorp.armeria.client.RpcClient;
30+
import com.linecorp.armeria.client.RpcPreprocessor;
31+
import com.linecorp.armeria.common.HttpRequest;
32+
import com.linecorp.armeria.common.HttpResponse;
33+
import com.linecorp.armeria.common.Request;
34+
import com.linecorp.armeria.common.Response;
35+
import com.linecorp.armeria.common.RpcRequest;
36+
import com.linecorp.armeria.common.RpcResponse;
37+
import com.linecorp.armeria.common.SessionProtocol;
38+
import com.linecorp.armeria.common.annotation.Nullable;
39+
import com.linecorp.armeria.internal.client.ClientRequestContextExtension;
40+
import com.linecorp.armeria.xds.client.endpoint.XdsEndpointGroup;
41+
import com.linecorp.armeria.xds.client.endpoint.XdsLoadBalancer;
42+
import com.linecorp.armeria.xds.internal.DelegatingHttpClient;
43+
import com.linecorp.armeria.xds.internal.DelegatingRpcClient;
44+
import com.linecorp.armeria.xds.internal.XdsCommonUtil;
45+
46+
/**
47+
* A factory which injects cluster-related filters.
48+
* <pre>{@code
49+
* [downstream preprocessors]
50+
* -> [router preprocessor]
51+
* -> [cluster preprocessor]
52+
* -> [endpoint selection]
53+
* -> [retry decorator]
54+
* -> [cluster decorator]
55+
* -> [upstream decorators]
56+
* }</pre>
57+
*/
58+
final class ClusterFilterFactory {
59+
60+
static final ClientDecoration DECORATION =
61+
ClientDecoration.builder()
62+
.add(ClusterFilterFactory::applyHttpClusterSettings)
63+
.addRpc(ClusterFilterFactory::applyRpcClusterSettings)
64+
.build();
65+
66+
private static final HttpClient CLUSTER_ONLY_HTTP_CLIENT =
67+
DECORATION.decorate(DelegatingHttpClient.of());
68+
private static final RpcClient CLUSTER_ONLY_RPC_CLIENT =
69+
DECORATION.rpcDecorate(DelegatingRpcClient.of());
70+
71+
private final XdsEndpointGroup endpointGroup;
72+
private final SessionProtocol sessionProtocol;
73+
74+
ClusterFilterFactory(XdsLoadBalancer loadBalancer,
75+
TransportSocketSnapshot transportSocket) {
76+
endpointGroup = XdsEndpointGroup.of(loadBalancer);
77+
sessionProtocol = transportSocket.clientTlsSpec() != null ?
78+
SessionProtocol.HTTPS : SessionProtocol.HTTP;
79+
}
80+
81+
HttpPreprocessor httpPreprocessor() {
82+
return this::execute;
83+
}
84+
85+
RpcPreprocessor rpcPreprocessor() {
86+
return this::execute;
87+
}
88+
89+
private <I extends Request, O extends Response> O execute(
90+
PreClient<I, O> delegate, PreClientRequestContext ctx, I req) throws Exception {
91+
ctx.setEndpointGroup(endpointGroup);
92+
ctx.setSessionProtocol(sessionProtocol);
93+
94+
final RouteEntry route = ctx.attr(XdsCommonUtil.SELECTED_ROUTE);
95+
final ClientRequestContextExtension ctxExt = ctx.as(ClientRequestContextExtension.class);
96+
if (ctxExt != null) {
97+
final HttpClient httpClient = route != null ?
98+
route.httpClient() : CLUSTER_ONLY_HTTP_CLIENT;
99+
final RpcClient rpcClient = route != null ?
100+
route.rpcClient() : CLUSTER_ONLY_RPC_CLIENT;
101+
ctxExt.httpClientCustomizer(actualClient -> {
102+
DelegatingHttpClient.setDelegate(ctx, actualClient);
103+
return httpClient;
104+
});
105+
ctxExt.rpcClientCustomizer(actualClient -> {
106+
DelegatingRpcClient.setDelegate(ctx, actualClient);
107+
return rpcClient;
108+
});
109+
}
110+
return delegate.execute(ctx, req);
111+
}
112+
113+
// Decorator logic — applies per-endpoint cluster settings (TLS)
114+
115+
private static HttpResponse applyHttpClusterSettings(
116+
HttpClient delegate, ClientRequestContext ctx, HttpRequest req) throws Exception {
117+
applyClusterSettings(ctx);
118+
return delegate.execute(ctx, req);
119+
}
120+
121+
private static RpcResponse applyRpcClusterSettings(
122+
RpcClient delegate, ClientRequestContext ctx, RpcRequest req) throws Exception {
123+
applyClusterSettings(ctx);
124+
return delegate.execute(ctx, req);
125+
}
126+
127+
private static void applyClusterSettings(ClientRequestContext ctx) {
128+
final Endpoint endpoint = ctx.endpoint();
129+
if (endpoint == null) {
130+
return;
131+
}
132+
final TransportSocketSnapshot transportSocket =
133+
endpoint.attr(XdsCommonUtil.TRANSPORT_SOCKET_SNAPSHOT_KEY);
134+
if (transportSocket == null) {
135+
return;
136+
}
137+
@Nullable
138+
ClientTlsSpec clientTlsSpec = transportSocket.clientTlsSpec();
139+
if (clientTlsSpec == null) {
140+
return;
141+
}
142+
final Set<String> alpnOverride = ctx.attr(XdsCommonUtil.ALPN_OVERRIDE_KEY);
143+
if (alpnOverride != null && !alpnOverride.isEmpty()) {
144+
clientTlsSpec = clientTlsSpec.toBuilder().alpnProtocols(alpnOverride).build();
145+
}
146+
ctx.setClientTlsSpec(clientTlsSpec);
147+
}
148+
}

0 commit comments

Comments
 (0)