From 0f968cc2161f26f604d5467513d61336207536c6 Mon Sep 17 00:00:00 2001 From: jrhee17 Date: Fri, 14 Feb 2025 18:45:00 +0900 Subject: [PATCH 1/2] introduce virtual host snapshot --- .../armeria/xds/ClusterResourceNode.java | 27 ++-- .../com/linecorp/armeria/xds/ClusterRoot.java | 2 +- .../linecorp/armeria/xds/ClusterSnapshot.java | 45 +----- .../armeria/xds/CompositeXdsStream.java | 2 +- .../armeria/xds/ConfigSourceMapper.java | 28 +--- .../linecorp/armeria/xds/ListenerRoot.java | 2 +- .../com/linecorp/armeria/xds/RouteEntry.java | 83 ++++++++++ .../armeria/xds/RouteResourceNode.java | 60 ++----- .../linecorp/armeria/xds/RouteSnapshot.java | 42 ++--- .../linecorp/armeria/xds/SotwXdsStream.java | 3 +- .../armeria/xds/StaticResourceUtils.java | 22 ++- .../armeria/xds/VirtualHostResourceNode.java | 122 +++++++++++++++ .../armeria/xds/VirtualHostSnapshot.java | 96 ++++++++++++ .../armeria/xds/VirtualHostXdsResource.java | 67 ++++++++ .../com/linecorp/armeria/xds/XdsType.java | 13 +- .../xds/client/endpoint/ClusterEntry.java | 21 ++- .../xds/client/endpoint/ClusterManager.java | 38 +++-- .../client/endpoint/DefaultLoadBalancer.java | 16 +- .../xds/client/endpoint/MetadataUtil.java | 45 +----- .../client/endpoint/SubsetLoadBalancer.java | 146 ++++++++++++------ .../xds/client/endpoint/XdsAttributeKeys.java | 3 + .../xds/client/endpoint/XdsLoadBalancer.java | 9 -- .../armeria/xds/AggregatingNodeTest.java | 38 +++-- .../armeria/xds/DynamicResourcesTest.java | 4 +- .../xds/MostlyStaticWithDynamicEdsTest.java | 4 +- .../armeria/xds/MultiConfigSourceTest.java | 21 ++- .../endpoint/RouteMetadataSubsetTest.java | 16 -- .../xds/client/endpoint/ZoneAwareTest.java | 8 - 28 files changed, 638 insertions(+), 345 deletions(-) create mode 100644 xds/src/main/java/com/linecorp/armeria/xds/RouteEntry.java create mode 100644 xds/src/main/java/com/linecorp/armeria/xds/VirtualHostResourceNode.java create mode 100644 xds/src/main/java/com/linecorp/armeria/xds/VirtualHostSnapshot.java create mode 100644 xds/src/main/java/com/linecorp/armeria/xds/VirtualHostXdsResource.java diff --git a/xds/src/main/java/com/linecorp/armeria/xds/ClusterResourceNode.java b/xds/src/main/java/com/linecorp/armeria/xds/ClusterResourceNode.java index cf0598ced8e..fb8d69835b4 100644 --- a/xds/src/main/java/com/linecorp/armeria/xds/ClusterResourceNode.java +++ b/xds/src/main/java/com/linecorp/armeria/xds/ClusterResourceNode.java @@ -18,7 +18,6 @@ import static com.google.common.base.Strings.isNullOrEmpty; import static com.linecorp.armeria.xds.XdsType.CLUSTER; -import static java.util.Objects.requireNonNull; import java.util.Objects; @@ -28,16 +27,10 @@ import io.envoyproxy.envoy.config.cluster.v3.Cluster.EdsClusterConfig; import io.envoyproxy.envoy.config.core.v3.ConfigSource; import io.envoyproxy.envoy.config.endpoint.v3.ClusterLoadAssignment; -import io.envoyproxy.envoy.config.route.v3.Route; -import io.envoyproxy.envoy.config.route.v3.VirtualHost; import io.grpc.Status; final class ClusterResourceNode extends AbstractResourceNodeWithPrimer { - @Nullable - private final VirtualHost virtualHost; - @Nullable - private final Route route; private final int index; private final EndpointSnapshotWatcher snapshotWatcher = new EndpointSnapshotWatcher(); private final SnapshotWatcher parentWatcher; @@ -49,19 +42,15 @@ final class ClusterResourceNode extends AbstractResourceNodeWithPrimer parentWatcher, - VirtualHost virtualHost, Route route, int index, ResourceNodeType resourceNodeType) { + @Nullable VirtualHostXdsResource primer, SnapshotWatcher parentWatcher, + int index, ResourceNodeType resourceNodeType) { super(xdsBootstrap, configSource, CLUSTER, resourceName, primer, parentWatcher, resourceNodeType); this.parentWatcher = parentWatcher; - this.virtualHost = requireNonNull(virtualHost, "virtualHost"); - this.route = requireNonNull(route, "route"); this.index = index; } @@ -86,10 +75,16 @@ public void doOnChanged(ClusterXdsResource resource) { children().add(node); xdsBootstrap().subscribe(node); } else { - parentWatcher.snapshotUpdated(new ClusterSnapshot(resource)); + final ClusterSnapshot clusterSnapshot = new ClusterSnapshot(resource); + parentWatcher.snapshotUpdated(clusterSnapshot); } } + @Override + public void close() { + super.close(); + } + private class EndpointSnapshotWatcher implements SnapshotWatcher { @Override public void snapshotUpdated(EndpointSnapshot newSnapshot) { @@ -100,8 +95,8 @@ public void snapshotUpdated(EndpointSnapshot newSnapshot) { if (!Objects.equals(newSnapshot.xdsResource().primer(), current)) { return; } - parentWatcher.snapshotUpdated( - new ClusterSnapshot(current, newSnapshot, virtualHost, route, index)); + final ClusterSnapshot clusterSnapshot = new ClusterSnapshot(current, newSnapshot, index); + parentWatcher.snapshotUpdated(clusterSnapshot); } @Override diff --git a/xds/src/main/java/com/linecorp/armeria/xds/ClusterRoot.java b/xds/src/main/java/com/linecorp/armeria/xds/ClusterRoot.java index 902088c72e2..c73b797a769 100644 --- a/xds/src/main/java/com/linecorp/armeria/xds/ClusterRoot.java +++ b/xds/src/main/java/com/linecorp/armeria/xds/ClusterRoot.java @@ -39,7 +39,7 @@ public final class ClusterRoot extends AbstractRoot { if (cluster != null) { node = staticCluster(xdsBootstrap, resourceName, this, cluster); } else { - final ConfigSource configSource = configSourceMapper.cdsConfigSource(null, resourceName); + final ConfigSource configSource = configSourceMapper.cdsConfigSource(resourceName); node = new ClusterResourceNode(configSource, resourceName, xdsBootstrap, null, this, ResourceNodeType.DYNAMIC); xdsBootstrap.subscribe(node); diff --git a/xds/src/main/java/com/linecorp/armeria/xds/ClusterSnapshot.java b/xds/src/main/java/com/linecorp/armeria/xds/ClusterSnapshot.java index 864c8f6870a..135cb9cc036 100644 --- a/xds/src/main/java/com/linecorp/armeria/xds/ClusterSnapshot.java +++ b/xds/src/main/java/com/linecorp/armeria/xds/ClusterSnapshot.java @@ -23,8 +23,6 @@ import com.linecorp.armeria.common.annotation.UnstableApi; import io.envoyproxy.envoy.config.cluster.v3.Cluster; -import io.envoyproxy.envoy.config.route.v3.Route; -import io.envoyproxy.envoy.config.route.v3.VirtualHost; /** * A snapshot of a {@link Cluster} resource. @@ -34,27 +32,17 @@ public final class ClusterSnapshot implements Snapshot { private final ClusterXdsResource clusterXdsResource; @Nullable private final EndpointSnapshot endpointSnapshot; - @Nullable - private final VirtualHost virtualHost; - @Nullable - private final Route route; private final int index; - ClusterSnapshot(ClusterXdsResource clusterXdsResource, EndpointSnapshot endpointSnapshot, - @Nullable VirtualHost virtualHost, @Nullable Route route, int index) { + ClusterSnapshot(ClusterXdsResource clusterXdsResource, + @Nullable EndpointSnapshot endpointSnapshot, int index) { this.clusterXdsResource = clusterXdsResource; this.endpointSnapshot = endpointSnapshot; - this.virtualHost = virtualHost; - this.route = route; this.index = index; } ClusterSnapshot(ClusterXdsResource clusterXdsResource) { - this.clusterXdsResource = clusterXdsResource; - endpointSnapshot = null; - virtualHost = null; - route = null; - index = -1; + this(clusterXdsResource, null, -1); } @Override @@ -70,22 +58,6 @@ public EndpointSnapshot endpointSnapshot() { return endpointSnapshot; } - /** - * The {@link VirtualHost} this {@link Cluster} belongs to. - */ - @Nullable - public VirtualHost virtualHost() { - return virtualHost; - } - - /** - * The {@link Route} this {@link Cluster} belongs to. - */ - @Nullable - public Route route() { - return route; - } - int index() { return index; } @@ -99,15 +71,13 @@ public boolean equals(Object object) { return false; } final ClusterSnapshot that = (ClusterSnapshot) object; - return index == that.index && Objects.equal(clusterXdsResource, that.clusterXdsResource) && - Objects.equal(endpointSnapshot, that.endpointSnapshot) && - Objects.equal(virtualHost, that.virtualHost) && - Objects.equal(route, that.route); + return Objects.equal(clusterXdsResource, that.clusterXdsResource) && + Objects.equal(endpointSnapshot, that.endpointSnapshot); } @Override public int hashCode() { - return Objects.hashCode(clusterXdsResource, endpointSnapshot, virtualHost, route, index); + return Objects.hashCode(clusterXdsResource, endpointSnapshot); } @Override @@ -116,9 +86,6 @@ public String toString() { .omitNullValues() .add("clusterXdsResource", clusterXdsResource) .add("endpointSnapshot", endpointSnapshot) - .add("virtualHost", virtualHost) - .add("route", route) - .add("index", index) .toString(); } } diff --git a/xds/src/main/java/com/linecorp/armeria/xds/CompositeXdsStream.java b/xds/src/main/java/com/linecorp/armeria/xds/CompositeXdsStream.java index 720fbf9c58b..1920c1ef898 100644 --- a/xds/src/main/java/com/linecorp/armeria/xds/CompositeXdsStream.java +++ b/xds/src/main/java/com/linecorp/armeria/xds/CompositeXdsStream.java @@ -34,7 +34,7 @@ final class CompositeXdsStream implements XdsStream { CompositeXdsStream(GrpcClientBuilder clientBuilder, Node node, Backoff backoff, EventExecutor eventLoop, XdsResponseHandler handler, SubscriberStorage subscriberStorage) { - for (XdsType type: XdsType.values()) { + for (XdsType type: XdsType.discoverableTypes()) { final SotwXdsStream stream = new SotwXdsStream( SotwDiscoveryStub.basic(type, clientBuilder), node, backoff, eventLoop, handler, subscriberStorage, EnumSet.of(type)); diff --git a/xds/src/main/java/com/linecorp/armeria/xds/ConfigSourceMapper.java b/xds/src/main/java/com/linecorp/armeria/xds/ConfigSourceMapper.java index 55bf07f141c..c00ace3a628 100644 --- a/xds/src/main/java/com/linecorp/armeria/xds/ConfigSourceMapper.java +++ b/xds/src/main/java/com/linecorp/armeria/xds/ConfigSourceMapper.java @@ -32,7 +32,7 @@ final class ConfigSourceMapper { @Nullable private final ConfigSource bootstrapAdsConfig; @Nullable - private ConfigSource parentConfigSource; + private final ConfigSource parentConfigSource; ConfigSourceMapper(Bootstrap bootstrap) { this(bootstrap, null); @@ -78,18 +78,7 @@ ConfigSource edsConfigSource(@Nullable ConfigSource configSource, String resourc throw new IllegalArgumentException("Cannot find an EDS config source for " + resourceName); } - ConfigSource cdsConfigSource(@Nullable ConfigSource configSource, String resourceName) { - if (configSource != null) { - if (configSource.hasApiConfigSource()) { - return configSource; - } - if (configSource.hasSelf() && parentConfigSource != null) { - return parentConfigSource; - } - if (configSource.hasAds() && bootstrapAdsConfig != null) { - return bootstrapAdsConfig; - } - } + ConfigSource cdsConfigSource(String resourceName) { if (bootstrapCdsConfig != null && bootstrapCdsConfig.hasApiConfigSource()) { return bootstrapCdsConfig; } @@ -114,18 +103,7 @@ ConfigSource rdsConfigSource(@Nullable ConfigSource configSource, String resourc throw new IllegalArgumentException("Cannot find an RDS config source for route: " + resourceName); } - ConfigSource ldsConfigSource(@Nullable ConfigSource configSource, String resourceName) { - if (configSource != null) { - if (configSource.hasApiConfigSource()) { - return configSource; - } - if (configSource.hasSelf() && parentConfigSource != null) { - return parentConfigSource; - } - if (configSource.hasAds() && bootstrapAdsConfig != null) { - return bootstrapAdsConfig; - } - } + ConfigSource ldsConfigSource(String resourceName) { if (bootstrapLdsConfig != null && bootstrapLdsConfig.hasApiConfigSource()) { return bootstrapLdsConfig; } diff --git a/xds/src/main/java/com/linecorp/armeria/xds/ListenerRoot.java b/xds/src/main/java/com/linecorp/armeria/xds/ListenerRoot.java index 97d53fb40e5..899fbacbd5f 100644 --- a/xds/src/main/java/com/linecorp/armeria/xds/ListenerRoot.java +++ b/xds/src/main/java/com/linecorp/armeria/xds/ListenerRoot.java @@ -39,7 +39,7 @@ public final class ListenerRoot extends AbstractRoot { node = new ListenerResourceNode(null, resourceName, xdsBootstrap, this, ResourceNodeType.STATIC); node.onChanged(listenerXdsResource); } else { - final ConfigSource configSource = configSourceMapper.ldsConfigSource(null, resourceName); + final ConfigSource configSource = configSourceMapper.ldsConfigSource(resourceName); node = new ListenerResourceNode(configSource, resourceName, xdsBootstrap, this, ResourceNodeType.DYNAMIC); xdsBootstrap.subscribe(node); diff --git a/xds/src/main/java/com/linecorp/armeria/xds/RouteEntry.java b/xds/src/main/java/com/linecorp/armeria/xds/RouteEntry.java new file mode 100644 index 00000000000..892abc52704 --- /dev/null +++ b/xds/src/main/java/com/linecorp/armeria/xds/RouteEntry.java @@ -0,0 +1,83 @@ +/* + * Copyright 2025 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.xds; + +import java.util.Objects; + +import com.google.common.base.MoreObjects; + +import com.linecorp.armeria.common.annotation.Nullable; + +import io.envoyproxy.envoy.config.route.v3.Route; +import io.envoyproxy.envoy.config.route.v3.RouteAction; + +/** + * Represents a {@link Route}. + */ +public final class RouteEntry { + + private final Route route; + @Nullable + private final ClusterSnapshot clusterSnapshot; + + RouteEntry(Route route, @Nullable ClusterSnapshot clusterSnapshot) { + this.route = route; + this.clusterSnapshot = clusterSnapshot; + } + + /** + * The {@link Route}. + */ + public Route route() { + return route; + } + + /** + * The {@link ClusterSnapshot} that is represented by {@link RouteAction#getCluster()}. + * If the {@link RouteAction} does not reference a cluster, the returned value may be {@code null}. + */ + @Nullable + public ClusterSnapshot clusterSnapshot() { + return clusterSnapshot; + } + + @Override + public boolean equals(Object o) { + if (this == o) { + return true; + } + if (o == null || getClass() != o.getClass()) { + return false; + } + final RouteEntry that = (RouteEntry) o; + return Objects.equals(route, that.route) && + Objects.equals(clusterSnapshot, that.clusterSnapshot); + } + + @Override + public int hashCode() { + return Objects.hash(route, clusterSnapshot); + } + + @Override + public String toString() { + return MoreObjects.toStringHelper(this) + .add("route", route) + .add("clusterSnapshot", clusterSnapshot) + .toString(); + } +} diff --git a/xds/src/main/java/com/linecorp/armeria/xds/RouteResourceNode.java b/xds/src/main/java/com/linecorp/armeria/xds/RouteResourceNode.java index a8faf7b2e88..35f325095c2 100644 --- a/xds/src/main/java/com/linecorp/armeria/xds/RouteResourceNode.java +++ b/xds/src/main/java/com/linecorp/armeria/xds/RouteResourceNode.java @@ -16,7 +16,6 @@ package com.linecorp.armeria.xds; -import static com.linecorp.armeria.xds.StaticResourceUtils.staticCluster; import static com.linecorp.armeria.xds.XdsType.ROUTE; import java.util.ArrayList; @@ -30,22 +29,17 @@ import com.linecorp.armeria.common.annotation.Nullable; -import io.envoyproxy.envoy.config.cluster.v3.Cluster; import io.envoyproxy.envoy.config.core.v3.ConfigSource; -import io.envoyproxy.envoy.config.route.v3.Route; -import io.envoyproxy.envoy.config.route.v3.Route.ActionCase; -import io.envoyproxy.envoy.config.route.v3.RouteAction; import io.envoyproxy.envoy.config.route.v3.RouteConfiguration; import io.envoyproxy.envoy.config.route.v3.VirtualHost; import io.grpc.Status; final class RouteResourceNode extends AbstractResourceNodeWithPrimer { - private final List clusterSnapshotList = new ArrayList<>(); - private final Set pending = new HashSet<>(); - private final ClusterSnapshotWatcher snapshotWatcher = new ClusterSnapshotWatcher(); + private final List virtualHostSnapshots = new ArrayList<>(); private final SnapshotWatcher parentWatcher; + private final VirtualHostSnapshotWatcher snapshotWatcher = new VirtualHostSnapshotWatcher(); RouteResourceNode(@Nullable ConfigSource configSource, String resourceName, XdsBootstrapImpl xdsBootstrap, @Nullable ListenerXdsResource primer, @@ -56,49 +50,27 @@ final class RouteResourceNode extends AbstractResourceNodeWithPrimer { + private class VirtualHostSnapshotWatcher implements SnapshotWatcher { @Override - public void snapshotUpdated(ClusterSnapshot newSnapshot) { + public void snapshotUpdated(VirtualHostSnapshot newSnapshot) { final RouteXdsResource current = currentResource(); if (current == null) { return; @@ -106,14 +78,14 @@ public void snapshotUpdated(ClusterSnapshot newSnapshot) { if (!Objects.equals(current, newSnapshot.xdsResource().primer())) { return; } - clusterSnapshotList.set(newSnapshot.index(), newSnapshot); + virtualHostSnapshots.set(newSnapshot.index(), newSnapshot); pending.remove(newSnapshot.index()); // checks if all clusters for the route have reported a snapshot if (!pending.isEmpty()) { return; } parentWatcher.snapshotUpdated( - new RouteSnapshot(current, ImmutableList.copyOf(clusterSnapshotList))); + new RouteSnapshot(current, ImmutableList.copyOf(virtualHostSnapshots))); } @Override diff --git a/xds/src/main/java/com/linecorp/armeria/xds/RouteSnapshot.java b/xds/src/main/java/com/linecorp/armeria/xds/RouteSnapshot.java index 448b33b911d..ae870117b37 100644 --- a/xds/src/main/java/com/linecorp/armeria/xds/RouteSnapshot.java +++ b/xds/src/main/java/com/linecorp/armeria/xds/RouteSnapshot.java @@ -16,11 +16,7 @@ package com.linecorp.armeria.xds; -import java.util.ArrayList; -import java.util.Collections; -import java.util.LinkedHashMap; import java.util.List; -import java.util.Map; import com.google.common.base.MoreObjects; import com.google.common.base.Objects; @@ -28,7 +24,6 @@ import com.linecorp.armeria.common.annotation.UnstableApi; import io.envoyproxy.envoy.config.route.v3.RouteConfiguration; -import io.envoyproxy.envoy.config.route.v3.VirtualHost; /** * A snapshot of a {@link RouteConfiguration} resource. @@ -37,22 +32,11 @@ public final class RouteSnapshot implements Snapshot { private final RouteXdsResource routeXdsResource; - private final List clusterSnapshots; + private final List virtualHostSnapshots; - private final Map> virtualHostMap; - - RouteSnapshot(RouteXdsResource routeXdsResource, List clusterSnapshots) { + RouteSnapshot(RouteXdsResource routeXdsResource, List virtualHostSnapshots) { this.routeXdsResource = routeXdsResource; - this.clusterSnapshots = clusterSnapshots; - - final LinkedHashMap> virtualHostMap = new LinkedHashMap<>(); - for (ClusterSnapshot clusterSnapshot: clusterSnapshots) { - final VirtualHost virtualHost = clusterSnapshot.virtualHost(); - assert virtualHost != null; - virtualHostMap.computeIfAbsent(virtualHost, ignored -> new ArrayList<>()) - .add(clusterSnapshot); - } - this.virtualHostMap = Collections.unmodifiableMap(virtualHostMap); + this.virtualHostSnapshots = virtualHostSnapshots; } @Override @@ -61,18 +45,10 @@ public RouteXdsResource xdsResource() { } /** - * A list of {@link ClusterSnapshot}s which belong to this {@link RouteConfiguration}. - */ - public List clusterSnapshots() { - return clusterSnapshots; - } - - /** - * A map of {@link VirtualHost}s to {@link ClusterSnapshot}s which belong - * to this {@link RouteConfiguration}. + * The virtual hosts represented by {@link RouteConfiguration#getVirtualHostsList()}. */ - public Map> virtualHostMap() { - return virtualHostMap; + public List virtualHostSnapshots() { + return virtualHostSnapshots; } @Override @@ -85,12 +61,12 @@ public boolean equals(Object object) { } final RouteSnapshot that = (RouteSnapshot) object; return Objects.equal(routeXdsResource, that.routeXdsResource) && - Objects.equal(clusterSnapshots, that.clusterSnapshots); + Objects.equal(virtualHostSnapshots, that.virtualHostSnapshots); } @Override public int hashCode() { - return Objects.hashCode(routeXdsResource, clusterSnapshots); + return Objects.hashCode(routeXdsResource, virtualHostSnapshots); } @Override @@ -98,7 +74,7 @@ public String toString() { return MoreObjects.toStringHelper(this) .omitNullValues() .add("routeXdsResource", routeXdsResource) - .add("clusterSnapshots", clusterSnapshots) + .add("virtualHostSnapshots", virtualHostSnapshots) .toString(); } } diff --git a/xds/src/main/java/com/linecorp/armeria/xds/SotwXdsStream.java b/xds/src/main/java/com/linecorp/armeria/xds/SotwXdsStream.java index 6ef26f321f7..46848230985 100644 --- a/xds/src/main/java/com/linecorp/armeria/xds/SotwXdsStream.java +++ b/xds/src/main/java/com/linecorp/armeria/xds/SotwXdsStream.java @@ -21,7 +21,6 @@ import java.util.Collection; import java.util.EnumMap; -import java.util.EnumSet; import java.util.Map; import java.util.Set; import java.util.concurrent.TimeUnit; @@ -74,7 +73,7 @@ final class SotwXdsStream implements XdsStream { XdsResponseHandler responseHandler, SubscriberStorage subscriberStorage) { this(stub, node, backoff, eventLoop, responseHandler, subscriberStorage, - EnumSet.allOf(XdsType.class)); + XdsType.discoverableTypes()); } SotwXdsStream(SotwDiscoveryStub stub, diff --git a/xds/src/main/java/com/linecorp/armeria/xds/StaticResourceUtils.java b/xds/src/main/java/com/linecorp/armeria/xds/StaticResourceUtils.java index 9f4cdc79ac7..5cc7643f3c3 100644 --- a/xds/src/main/java/com/linecorp/armeria/xds/StaticResourceUtils.java +++ b/xds/src/main/java/com/linecorp/armeria/xds/StaticResourceUtils.java @@ -20,7 +20,6 @@ import io.envoyproxy.envoy.config.cluster.v3.Cluster; import io.envoyproxy.envoy.config.endpoint.v3.ClusterLoadAssignment; -import io.envoyproxy.envoy.config.route.v3.Route; import io.envoyproxy.envoy.config.route.v3.RouteConfiguration; import io.envoyproxy.envoy.config.route.v3.VirtualHost; @@ -41,6 +40,19 @@ static RouteResourceNode staticRoute(XdsBootstrapImpl xdsBootstrap, String resou return node; } + static VirtualHostResourceNode staticVirtualHost( + XdsBootstrapImpl xdsBootstrap, String resourceName, + RouteXdsResource primer, + SnapshotWatcher parentWatcher, + int index, VirtualHost virtualHost) { + final VirtualHostResourceNode node = + new VirtualHostResourceNode(null, resourceName, xdsBootstrap, + primer, parentWatcher, index, STATIC); + final VirtualHostXdsResource resource = new VirtualHostXdsResource(virtualHost); + node.onChanged(resource); + return node; + } + static ClusterResourceNode staticCluster(XdsBootstrapImpl xdsBootstrap, String resourceName, SnapshotWatcher parentWatcher, Cluster cluster) { @@ -51,13 +63,11 @@ static ClusterResourceNode staticCluster(XdsBootstrapImpl xdsBootstrap, String r } static ClusterResourceNode staticCluster(XdsBootstrapImpl xdsBootstrap, String resourceName, - RouteXdsResource primer, + VirtualHostXdsResource primer, SnapshotWatcher parentWatcher, - VirtualHost virtualHost, Route route, int index, - Cluster cluster) { + int index, Cluster cluster) { final ClusterResourceNode node = new ClusterResourceNode(null, resourceName, xdsBootstrap, - primer, parentWatcher, virtualHost, route, - index, STATIC); + primer, parentWatcher, index, STATIC); setClusterXdsResourceToNode(cluster, node); return node; } diff --git a/xds/src/main/java/com/linecorp/armeria/xds/VirtualHostResourceNode.java b/xds/src/main/java/com/linecorp/armeria/xds/VirtualHostResourceNode.java new file mode 100644 index 00000000000..20c2e6c4c56 --- /dev/null +++ b/xds/src/main/java/com/linecorp/armeria/xds/VirtualHostResourceNode.java @@ -0,0 +1,122 @@ +/* + * Copyright 2025 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.xds; + +import static com.linecorp.armeria.xds.StaticResourceUtils.staticCluster; + +import java.util.ArrayList; +import java.util.HashSet; +import java.util.List; +import java.util.Objects; +import java.util.Set; + +import com.google.common.collect.ImmutableList; + +import com.linecorp.armeria.common.annotation.Nullable; + +import io.envoyproxy.envoy.config.cluster.v3.Cluster; +import io.envoyproxy.envoy.config.core.v3.ConfigSource; +import io.envoyproxy.envoy.config.route.v3.Route; +import io.envoyproxy.envoy.config.route.v3.Route.ActionCase; +import io.envoyproxy.envoy.config.route.v3.RouteAction; +import io.grpc.Status; + +final class VirtualHostResourceNode extends AbstractResourceNodeWithPrimer { + + private final Set pending = new HashSet<>(); + private final List clusterSnapshots = new ArrayList<>(); + private final ClusterSnapshotWatcher snapshotWatcher = new ClusterSnapshotWatcher(); + private final SnapshotWatcher parentWatcher; + private final int index; + + VirtualHostResourceNode(@Nullable ConfigSource configSource, String resourceName, + XdsBootstrapImpl xdsBootstrap, @Nullable RouteXdsResource primer, + SnapshotWatcher parentWatcher, int index, + ResourceNodeType resourceNodeType) { + super(xdsBootstrap, configSource, XdsType.VIRTUAL_HOST, resourceName, primer, parentWatcher, + resourceNodeType); + this.parentWatcher = parentWatcher; + this.index = index; + } + + @Override + void doOnChanged(VirtualHostXdsResource resource) { + pending.clear(); + for (Route route: resource.resource().getRoutesList()) { + final RouteAction routeAction = route.getRoute(); + final String clusterName = routeAction.getCluster(); + + // add a dummy element to the index list so that we can call List.set later + // without incurring an IndexOutOfBoundException when a snapshot is updated + clusterSnapshots.add(null); + + if (route.getActionCase() != ActionCase.ROUTE) { + continue; + } + + final int index = clusterSnapshots.size() - 1; + pending.add(index); + + final Cluster cluster = xdsBootstrap().bootstrapClusters().cluster(clusterName); + final ClusterResourceNode node; + if (cluster != null) { + node = staticCluster(xdsBootstrap(), clusterName, resource, snapshotWatcher, + index, cluster); + children().add(node); + } else { + final ConfigSource configSource = + configSourceMapper().cdsConfigSource(clusterName); + node = new ClusterResourceNode(configSource, clusterName, xdsBootstrap(), + resource, snapshotWatcher, index, ResourceNodeType.DYNAMIC); + children().add(node); + xdsBootstrap().subscribe(node); + } + } + } + + private class ClusterSnapshotWatcher implements SnapshotWatcher { + + @Override + public void snapshotUpdated(ClusterSnapshot newSnapshot) { + final VirtualHostXdsResource current = currentResource(); + if (current == null) { + return; + } + if (!Objects.equals(current, newSnapshot.xdsResource().primer())) { + return; + } + clusterSnapshots.set(newSnapshot.index(), newSnapshot); + pending.remove(newSnapshot.index()); + // checks if all clusters for the route have reported a snapshot + if (!pending.isEmpty()) { + return; + } + parentWatcher.snapshotUpdated( + new VirtualHostSnapshot(current, ImmutableList.copyOf(clusterSnapshots), index)); + } + + @Override + public void onError(XdsType type, Status status) { + parentWatcher.onError(type, status); + } + + @Override + public void onMissing(XdsType type, String resourceName) { + parentWatcher.onMissing(type, resourceName); + } + } +} diff --git a/xds/src/main/java/com/linecorp/armeria/xds/VirtualHostSnapshot.java b/xds/src/main/java/com/linecorp/armeria/xds/VirtualHostSnapshot.java new file mode 100644 index 00000000000..7af63d528b1 --- /dev/null +++ b/xds/src/main/java/com/linecorp/armeria/xds/VirtualHostSnapshot.java @@ -0,0 +1,96 @@ +/* + * Copyright 2025 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.xds; + +import java.util.List; +import java.util.Objects; + +import com.google.common.base.MoreObjects; +import com.google.common.collect.ImmutableList; + +import com.linecorp.armeria.common.annotation.Nullable; + +import io.envoyproxy.envoy.config.route.v3.Route; +import io.envoyproxy.envoy.config.route.v3.VirtualHost; + +/** + * A snapshot of a {@link VirtualHost}. + */ +public final class VirtualHostSnapshot implements Snapshot { + + private final VirtualHostXdsResource virtualHostXdsResource; + private final List routeEntries; + private final int index; + + VirtualHostSnapshot(VirtualHostXdsResource virtualHostXdsResource, + List<@Nullable ClusterSnapshot> clusterSnapshots, int index) { + this.virtualHostXdsResource = virtualHostXdsResource; + assert clusterSnapshots.size() == virtualHostXdsResource.resource().getRoutesCount(); + + final ImmutableList.Builder routeEntriesBuilder = ImmutableList.builder(); + for (int i = 0; i < clusterSnapshots.size(); i++) { + final ClusterSnapshot clusterSnapshot = clusterSnapshots.get(i); + final Route route = virtualHostXdsResource.resource().getRoutes(i); + routeEntriesBuilder.add(new RouteEntry(route, clusterSnapshot)); + } + routeEntries = routeEntriesBuilder.build(); + this.index = index; + } + + @Override + public VirtualHostXdsResource xdsResource() { + return virtualHostXdsResource; + } + + /** + * A list of routes corresponding to {@link VirtualHost#getRoutesList()}. + */ + public List routeEntries() { + return routeEntries; + } + + int index() { + return index; + } + + @Override + public boolean equals(Object o) { + if (this == o) { + return true; + } + if (o == null || getClass() != o.getClass()) { + return false; + } + final VirtualHostSnapshot that = (VirtualHostSnapshot) o; + return index == that.index && + Objects.equals(virtualHostXdsResource, that.virtualHostXdsResource) && + Objects.equals(routeEntries, that.routeEntries); + } + + @Override + public int hashCode() { + return Objects.hash(virtualHostXdsResource, routeEntries, index); + } + + @Override + public String toString() { + return MoreObjects.toStringHelper(this) + .add("virtualHostXdsResource", virtualHostXdsResource) + .add("routeEntries", routeEntries) + .toString(); + } +} diff --git a/xds/src/main/java/com/linecorp/armeria/xds/VirtualHostXdsResource.java b/xds/src/main/java/com/linecorp/armeria/xds/VirtualHostXdsResource.java new file mode 100644 index 00000000000..adeaf701e82 --- /dev/null +++ b/xds/src/main/java/com/linecorp/armeria/xds/VirtualHostXdsResource.java @@ -0,0 +1,67 @@ +/* + * Copyright 2025 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.xds; + +import com.linecorp.armeria.common.annotation.Nullable; + +import io.envoyproxy.envoy.config.route.v3.VirtualHost; + +/** + * A resource object for a {@link VirtualHost}. + */ +public final class VirtualHostXdsResource extends XdsResourceWithPrimer { + + private final VirtualHost virtualHost; + @Nullable + private final XdsResource primer; + + VirtualHostXdsResource(VirtualHost virtualHost) { + this.virtualHost = virtualHost; + primer = null; + } + + VirtualHostXdsResource(VirtualHost virtualHost, @Nullable XdsResource primer) { + this.virtualHost = virtualHost; + this.primer = primer; + } + + @Override + public XdsType type() { + return XdsType.VIRTUAL_HOST; + } + + @Override + public VirtualHost resource() { + return virtualHost; + } + + @Override + public String name() { + return virtualHost.getName(); + } + + @Override + VirtualHostXdsResource withPrimer(@Nullable XdsResource primer) { + return new VirtualHostXdsResource(virtualHost, primer); + } + + @Nullable + @Override + XdsResource primer() { + return primer; + } +} diff --git a/xds/src/main/java/com/linecorp/armeria/xds/XdsType.java b/xds/src/main/java/com/linecorp/armeria/xds/XdsType.java index da154403d2c..c5a9ec39143 100644 --- a/xds/src/main/java/com/linecorp/armeria/xds/XdsType.java +++ b/xds/src/main/java/com/linecorp/armeria/xds/XdsType.java @@ -16,6 +16,9 @@ package com.linecorp.armeria.xds; +import java.util.EnumSet; +import java.util.Set; + import com.linecorp.armeria.common.annotation.UnstableApi; /** @@ -23,10 +26,14 @@ */ @UnstableApi public enum XdsType { + LISTENER("type.googleapis.com/envoy.config.listener.v3.Listener"), ROUTE("type.googleapis.com/envoy.config.route.v3.RouteConfiguration"), CLUSTER("type.googleapis.com/envoy.config.cluster.v3.Cluster"), - ENDPOINT("type.googleapis.com/envoy.config.endpoint.v3.ClusterLoadAssignment"); + ENDPOINT("type.googleapis.com/envoy.config.endpoint.v3.ClusterLoadAssignment"), + VIRTUAL_HOST("type.googleapis.com/envoy.config.route.v3.VirtualHost"); + + private static final Set discoverableTypes = EnumSet.of(LISTENER, ROUTE, CLUSTER, ENDPOINT); private final String typeUrl; @@ -40,4 +47,8 @@ public enum XdsType { public String typeUrl() { return typeUrl; } + + static Set discoverableTypes() { + return discoverableTypes; + } } diff --git a/xds/src/main/java/com/linecorp/armeria/xds/client/endpoint/ClusterEntry.java b/xds/src/main/java/com/linecorp/armeria/xds/client/endpoint/ClusterEntry.java index 70c1e63f4fc..c84a271b619 100644 --- a/xds/src/main/java/com/linecorp/armeria/xds/client/endpoint/ClusterEntry.java +++ b/xds/src/main/java/com/linecorp/armeria/xds/client/endpoint/ClusterEntry.java @@ -34,8 +34,8 @@ import com.linecorp.armeria.common.util.AbstractListenable; import com.linecorp.armeria.common.util.AsyncCloseable; import com.linecorp.armeria.xds.ClusterSnapshot; +import com.linecorp.armeria.xds.RouteEntry; import com.linecorp.armeria.xds.client.endpoint.ClusterManager.LocalCluster; -import com.linecorp.armeria.xds.client.endpoint.LocalityRoutingStateFactory.LocalityRoutingState; import io.netty.util.concurrent.EventExecutor; @@ -53,13 +53,17 @@ final class ClusterEntry extends AbstractListenable implements private final EndpointsPool endpointsPool; @Nullable private final LocalCluster localCluster; + @Nullable + private final RouteEntry routeEntry; private final EventExecutor eventExecutor; private boolean closing; - ClusterEntry(EventExecutor eventExecutor, @Nullable LocalCluster localCluster) { + ClusterEntry(EventExecutor eventExecutor, @Nullable LocalCluster localCluster, + @Nullable RouteEntry routeEntry) { this.eventExecutor = eventExecutor; endpointsPool = new EndpointsPool(eventExecutor); this.localCluster = localCluster; + this.routeEntry = routeEntry; if (localCluster != null) { localCluster.clusterEntry().addListener(localClusterEntryListener, true); } @@ -107,16 +111,17 @@ void tryRefresh() { logger.trace("XdsEndpointGroup is using a new PrioritySet({})", prioritySet); } - LocalityRoutingState localityRoutingState = null; + PrioritySet localPrioritySet = null; + final XdsLoadBalancer localLoadBalancer = this.localLoadBalancer; if (localLoadBalancer != null) { assert localCluster != null; - localityRoutingState = localCluster.stateFactory().create(prioritySet, - localLoadBalancer.prioritySet()); - logger.trace("Local routing is enabled with LocalityRoutingState({})", localityRoutingState); + localPrioritySet = localLoadBalancer.prioritySet(); + logger.trace("Local routing is enabled with local prioritySet({})", localPrioritySet); } - XdsLoadBalancer loadBalancer = new DefaultLoadBalancer(prioritySet, localityRoutingState); + XdsLoadBalancer loadBalancer = new DefaultLoadBalancer(prioritySet, localCluster, localPrioritySet); if (clusterSnapshot.xdsResource().resource().hasLbSubsetConfig()) { - loadBalancer = new SubsetLoadBalancer(prioritySet, loadBalancer); + loadBalancer = new SubsetLoadBalancer(prioritySet, loadBalancer, localCluster, localPrioritySet, + routeEntry); } this.loadBalancer = loadBalancer; notifyListeners(loadBalancer); diff --git a/xds/src/main/java/com/linecorp/armeria/xds/client/endpoint/ClusterManager.java b/xds/src/main/java/com/linecorp/armeria/xds/client/endpoint/ClusterManager.java index 830466f6c5b..5b633595623 100644 --- a/xds/src/main/java/com/linecorp/armeria/xds/client/endpoint/ClusterManager.java +++ b/xds/src/main/java/com/linecorp/armeria/xds/client/endpoint/ClusterManager.java @@ -54,8 +54,10 @@ import com.linecorp.armeria.xds.ClusterSnapshot; import com.linecorp.armeria.xds.ListenerRoot; import com.linecorp.armeria.xds.ListenerSnapshot; +import com.linecorp.armeria.xds.RouteEntry; import com.linecorp.armeria.xds.RouteSnapshot; import com.linecorp.armeria.xds.SnapshotWatcher; +import com.linecorp.armeria.xds.VirtualHostSnapshot; import com.linecorp.armeria.xds.XdsBootstrap; import com.linecorp.armeria.xds.client.endpoint.ClusterManager.State; @@ -100,7 +102,7 @@ final class ClusterManager implements SnapshotWatcher, AsyncCl eventLoop = CommonPools.workerGroup().next(); listenerRoot = null; localCluster = null; - final ClusterEntry clusterEntry = new ClusterEntry(eventLoop, null); + final ClusterEntry clusterEntry = new ClusterEntry(eventLoop, null, null); clusterEntry.addListener(ignored -> notifyListeners(), true); clusterEntry.updateClusterSnapshot(clusterSnapshot); clusterEntries = @@ -123,26 +125,30 @@ public void snapshotUpdated(ListenerSnapshot listenerSnapshot) { if (closed) { return; } - final RouteSnapshot routeSnapshot = listenerSnapshot.routeSnapshot(); - final List clusterSnapshots = - routeSnapshot != null ? routeSnapshot.clusterSnapshots() : ImmutableList.of(); + final ClusterEntries clusterEntries = this.clusterEntries; final Map oldClusterEntries = clusterEntries.clusterEntriesMap; // ImmutableMap is used because it is important that the entries are added in order of // ClusterSnapshot#index so that the first matching route is selected in #selectNow final ImmutableMap.Builder mappingBuilder = ImmutableMap.builder(); - for (ClusterSnapshot clusterSnapshot : clusterSnapshots) { - if (clusterSnapshot.endpointSnapshot() == null) { - continue; - } - final String clusterName = clusterSnapshot.xdsResource().name(); - ClusterEntry clusterEntry = oldClusterEntries.get(clusterName); - if (clusterEntry == null) { - clusterEntry = new ClusterEntry(eventLoop, localCluster); - clusterEntry.addListener(ignored -> notifyListeners(), false); + final RouteSnapshot routeSnapshot = listenerSnapshot.routeSnapshot(); + final List virtualHostSnapshots = + routeSnapshot == null ? ImmutableList.of() : routeSnapshot.virtualHostSnapshots(); + for (VirtualHostSnapshot virtualHostSnapshot: virtualHostSnapshots) { + for (RouteEntry routeEntry: virtualHostSnapshot.routeEntries()) { + final ClusterSnapshot clusterSnapshot = routeEntry.clusterSnapshot(); + if (clusterSnapshot == null || clusterSnapshot.endpointSnapshot() == null) { + continue; + } + final String clusterName = clusterSnapshot.xdsResource().name(); + ClusterEntry clusterEntry = oldClusterEntries.get(clusterName); + if (clusterEntry == null) { + clusterEntry = new ClusterEntry(eventLoop, localCluster, routeEntry); + clusterEntry.addListener(ignored -> notifyListeners(), false); + } + clusterEntry.updateClusterSnapshot(clusterSnapshot); + mappingBuilder.put(clusterName, clusterEntry); } - clusterEntry.updateClusterSnapshot(clusterSnapshot); - mappingBuilder.put(clusterName, clusterEntry); } final ImmutableMap newClusterEntriesMap = mappingBuilder.build(); this.clusterEntries = new ClusterEntries(listenerSnapshot, newClusterEntriesMap); @@ -338,7 +344,7 @@ static final class LocalCluster implements AsyncCloseable { private LocalCluster(String localClusterName, XdsBootstrap xdsBootstrap) { final Node node = xdsBootstrap.bootstrap().getNode(); localityRoutingStateFactory = new LocalityRoutingStateFactory(node.getLocality()); - clusterEntry = new ClusterEntry(xdsBootstrap.eventLoop(), null); + clusterEntry = new ClusterEntry(xdsBootstrap.eventLoop(), null, null); localClusterRoot = xdsBootstrap.clusterRoot(localClusterName); localClusterRoot.addSnapshotWatcher(clusterEntry::updateClusterSnapshot); } diff --git a/xds/src/main/java/com/linecorp/armeria/xds/client/endpoint/DefaultLoadBalancer.java b/xds/src/main/java/com/linecorp/armeria/xds/client/endpoint/DefaultLoadBalancer.java index 381977955b7..9a8d5c97599 100644 --- a/xds/src/main/java/com/linecorp/armeria/xds/client/endpoint/DefaultLoadBalancer.java +++ b/xds/src/main/java/com/linecorp/armeria/xds/client/endpoint/DefaultLoadBalancer.java @@ -28,6 +28,7 @@ import com.linecorp.armeria.client.Endpoint; import com.linecorp.armeria.client.endpoint.EndpointGroup; import com.linecorp.armeria.common.annotation.Nullable; +import com.linecorp.armeria.xds.client.endpoint.ClusterManager.LocalCluster; import com.linecorp.armeria.xds.client.endpoint.DefaultLbStateFactory.DefaultLbState; import com.linecorp.armeria.xds.client.endpoint.LocalityRoutingStateFactory.LocalityRoutingState; import com.linecorp.armeria.xds.client.endpoint.LocalityRoutingStateFactory.State; @@ -41,10 +42,15 @@ final class DefaultLoadBalancer implements XdsLoadBalancer { @Nullable private final LocalityRoutingState localityRoutingState; - DefaultLoadBalancer(PrioritySet prioritySet, @Nullable LocalityRoutingState localityRoutingState) { + DefaultLoadBalancer(PrioritySet prioritySet, @Nullable LocalCluster localCluster, + @Nullable PrioritySet localPrioritySet) { lbState = DefaultLbStateFactory.newInstance(prioritySet); this.prioritySet = prioritySet; - this.localityRoutingState = localityRoutingState; + if (localCluster != null && localPrioritySet != null) { + localityRoutingState = localCluster.stateFactory().create(prioritySet, localPrioritySet); + } else { + localityRoutingState = null; + } } @Override @@ -185,12 +191,6 @@ public PrioritySet prioritySet() { return prioritySet; } - @Override - @Nullable - public LocalityRoutingState localityRoutingState() { - return localityRoutingState; - } - static class PriorityAndAvailability { final int priority; final HostAvailability hostAvailability; diff --git a/xds/src/main/java/com/linecorp/armeria/xds/client/endpoint/MetadataUtil.java b/xds/src/main/java/com/linecorp/armeria/xds/client/endpoint/MetadataUtil.java index 3cbe7ea0260..a014d647811 100644 --- a/xds/src/main/java/com/linecorp/armeria/xds/client/endpoint/MetadataUtil.java +++ b/xds/src/main/java/com/linecorp/armeria/xds/client/endpoint/MetadataUtil.java @@ -18,50 +18,21 @@ import static com.linecorp.armeria.xds.client.endpoint.XdsConstants.SUBSET_LOAD_BALANCING_FILTER_NAME; -import java.util.Map; - -import com.google.protobuf.ProtocolStringList; import com.google.protobuf.Struct; -import com.google.protobuf.Value; -import com.linecorp.armeria.xds.ClusterSnapshot; +import com.linecorp.armeria.client.ClientRequestContext; -import io.envoyproxy.envoy.config.cluster.v3.Cluster.LbSubsetConfig; -import io.envoyproxy.envoy.config.cluster.v3.Cluster.LbSubsetConfig.LbSubsetSelector; -import io.envoyproxy.envoy.config.route.v3.Route; -import io.envoyproxy.envoy.config.route.v3.RouteAction; +import io.envoyproxy.envoy.config.core.v3.Metadata; final class MetadataUtil { - static Struct filterMetadata(ClusterSnapshot clusterSnapshot) { - final Route route = clusterSnapshot.route(); - if (route == null) { - return Struct.getDefaultInstance(); - } - final RouteAction action = route.getRoute(); - return action.getMetadataMatch().getFilterMetadataOrDefault(SUBSET_LOAD_BALANCING_FILTER_NAME, - Struct.getDefaultInstance()); - } - - static boolean findMatchedSubsetSelector(LbSubsetConfig lbSubsetConfig, Struct filterMetadata) { - for (LbSubsetSelector subsetSelector : lbSubsetConfig.getSubsetSelectorsList()) { - final ProtocolStringList keysList = subsetSelector.getKeysList(); - if (filterMetadata.getFieldsCount() != keysList.size()) { - continue; - } - boolean found = true; - final Map filterMetadataMap = filterMetadata.getFieldsMap(); - for (String key : filterMetadataMap.keySet()) { - if (!keysList.contains(key)) { - found = false; - break; - } - } - if (found) { - return true; - } + static Struct filterMetadata(ClientRequestContext ctx, Struct defaultMetadata) { + final Metadata metadataMatch = ctx.attr(XdsAttributeKeys.ROUTE_METADATA_MATCH); + if (metadataMatch == null) { + return defaultMetadata; } - return false; + return metadataMatch.getFilterMetadataOrDefault(SUBSET_LOAD_BALANCING_FILTER_NAME, + defaultMetadata); } private MetadataUtil() {} diff --git a/xds/src/main/java/com/linecorp/armeria/xds/client/endpoint/SubsetLoadBalancer.java b/xds/src/main/java/com/linecorp/armeria/xds/client/endpoint/SubsetLoadBalancer.java index 502c9bb31db..8bd4ceee577 100644 --- a/xds/src/main/java/com/linecorp/armeria/xds/client/endpoint/SubsetLoadBalancer.java +++ b/xds/src/main/java/com/linecorp/armeria/xds/client/endpoint/SubsetLoadBalancer.java @@ -17,63 +17,80 @@ package com.linecorp.armeria.xds.client.endpoint; import static com.linecorp.armeria.xds.client.endpoint.MetadataUtil.filterMetadata; -import static com.linecorp.armeria.xds.client.endpoint.MetadataUtil.findMatchedSubsetSelector; -import static com.linecorp.armeria.xds.client.endpoint.XdsEndpointUtil.convertEndpoints; +import static com.linecorp.armeria.xds.client.endpoint.XdsConstants.SUBSET_LOAD_BALANCING_FILTER_NAME; +import java.util.ArrayList; +import java.util.HashMap; import java.util.List; +import java.util.Map; +import java.util.Map.Entry; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import com.google.common.base.MoreObjects; +import com.google.common.collect.ImmutableMap; +import com.google.protobuf.ProtocolStringList; import com.google.protobuf.Struct; import com.linecorp.armeria.client.ClientRequestContext; import com.linecorp.armeria.client.Endpoint; import com.linecorp.armeria.common.annotation.Nullable; import com.linecorp.armeria.xds.ClusterSnapshot; -import com.linecorp.armeria.xds.client.endpoint.LocalityRoutingStateFactory.LocalityRoutingState; +import com.linecorp.armeria.xds.RouteEntry; +import com.linecorp.armeria.xds.client.endpoint.ClusterManager.LocalCluster; import io.envoyproxy.envoy.config.cluster.v3.Cluster; import io.envoyproxy.envoy.config.cluster.v3.Cluster.LbSubsetConfig; import io.envoyproxy.envoy.config.cluster.v3.Cluster.LbSubsetConfig.LbSubsetFallbackPolicy; +import io.envoyproxy.envoy.config.cluster.v3.Cluster.LbSubsetConfig.LbSubsetSelector; +import io.envoyproxy.envoy.config.endpoint.v3.LbEndpoint; final class SubsetLoadBalancer implements XdsLoadBalancer { private static final Logger logger = LoggerFactory.getLogger(SubsetLoadBalancer.class); - private final LoadBalancer loadBalancer; private final PrioritySet prioritySet; + private final LoadBalancer allEndpointsLoadBalancer; @Nullable - private final LocalityRoutingState localityRoutingState; + private final LocalCluster localCluster; + @Nullable + private final PrioritySet localPrioritySet; + + private final Map subsetLoadBalancers; + private final LbSubsetConfig lbSubsetConfig; + private final LbSubsetFallbackPolicy fallbackPolicy; + + private final Struct defaultMetadataMatch; + + SubsetLoadBalancer(PrioritySet prioritySet, LoadBalancer allEndpointsLoadBalancer, + @Nullable LocalCluster localCluster, @Nullable PrioritySet localPrioritySet, + @Nullable RouteEntry routeEntry) { + this.allEndpointsLoadBalancer = allEndpointsLoadBalancer; + this.localCluster = localCluster; + this.localPrioritySet = localPrioritySet; - SubsetLoadBalancer(PrioritySet prioritySet, XdsLoadBalancer allEndpointsLoadBalancer) { - loadBalancer = createSubsetLoadBalancer(prioritySet, allEndpointsLoadBalancer); + final ClusterSnapshot clusterSnapshot = prioritySet.clusterSnapshot(); + final Cluster cluster = clusterSnapshot.xdsResource().resource(); + lbSubsetConfig = cluster.getLbSubsetConfig(); + fallbackPolicy = lbSubsetFallbackPolicy(lbSubsetConfig); + + subsetLoadBalancers = createSubsetLoadBalancers(prioritySet); this.prioritySet = prioritySet; - localityRoutingState = allEndpointsLoadBalancer.localityRoutingState(); - } - @Override - @Nullable - public Endpoint selectNow(ClientRequestContext ctx) { - return loadBalancer.selectNow(ctx); + defaultMetadataMatch = defaultMetadataMatch(routeEntry); } - private LoadBalancer createSubsetLoadBalancer(PrioritySet prioritySet, - LoadBalancer allEndpointsLoadBalancer) { - final ClusterSnapshot clusterSnapshot = prioritySet.clusterSnapshot(); - final Struct filterMetadata = filterMetadata(clusterSnapshot); - if (filterMetadata.getFieldsCount() == 0) { - // No metadata. Use the whole endpoints. - return allEndpointsLoadBalancer; + private static Struct defaultMetadataMatch(@Nullable RouteEntry routeEntry) { + if (routeEntry == null) { + return Struct.getDefaultInstance(); } + return routeEntry.route().getRoute().getMetadataMatch() + .getFilterMetadataOrDefault(SUBSET_LOAD_BALANCING_FILTER_NAME, + Struct.getDefaultInstance()); + } - final Cluster cluster = clusterSnapshot.xdsResource().resource(); - final LbSubsetConfig lbSubsetConfig = cluster.getLbSubsetConfig(); - if (lbSubsetConfig == LbSubsetConfig.getDefaultInstance()) { - // Route metadata exists but no lbSubsetConfig. Use NO_FALLBACK. - return NOOP; - } + private static LbSubsetFallbackPolicy lbSubsetFallbackPolicy(LbSubsetConfig lbSubsetConfig) { LbSubsetFallbackPolicy fallbackPolicy = lbSubsetConfig.getFallbackPolicy(); if (!(fallbackPolicy == LbSubsetFallbackPolicy.NO_FALLBACK || fallbackPolicy == LbSubsetFallbackPolicy.ANY_ENDPOINT)) { @@ -81,34 +98,69 @@ private LoadBalancer createSubsetLoadBalancer(PrioritySet prioritySet, fallbackPolicy, LbSubsetFallbackPolicy.NO_FALLBACK); fallbackPolicy = LbSubsetFallbackPolicy.NO_FALLBACK; } + return fallbackPolicy; + } - if (!findMatchedSubsetSelector(lbSubsetConfig, filterMetadata)) { - if (fallbackPolicy == LbSubsetFallbackPolicy.NO_FALLBACK) { - return NOOP; - } - return allEndpointsLoadBalancer; + @Override + @Nullable + public Endpoint selectNow(ClientRequestContext ctx) { + final Struct filterMetadata = filterMetadata(ctx, defaultMetadataMatch); + final LoadBalancer subsetLoadBalancer = subsetLoadBalancers.get(filterMetadata); + if (subsetLoadBalancer != null) { + return subsetLoadBalancer.selectNow(ctx); } - final List endpoints = convertEndpoints(prioritySet.endpoints(), - filterMetadata); - if (endpoints.isEmpty()) { - if (fallbackPolicy == LbSubsetFallbackPolicy.NO_FALLBACK) { - return NOOP; - } - return allEndpointsLoadBalancer; + if (fallbackPolicy == LbSubsetFallbackPolicy.NO_FALLBACK) { + return null; } - return createSubsetLoadBalancer(endpoints, clusterSnapshot); + assert fallbackPolicy == LbSubsetFallbackPolicy.ANY_ENDPOINT; + return allEndpointsLoadBalancer.selectNow(ctx); } - private LoadBalancer createSubsetLoadBalancer(List endpoints, - ClusterSnapshot clusterSnapshot) { - final PrioritySet subsetPrioritySet = new PriorityStateManager(clusterSnapshot, endpoints).build(); - return new DefaultLoadBalancer(subsetPrioritySet, localityRoutingState); + private Map createSubsetLoadBalancers(PrioritySet prioritySet) { + final ClusterSnapshot clusterSnapshot = prioritySet.clusterSnapshot(); + + final Map> endpointsPerFilterStruct = new HashMap<>(); + for (LbSubsetSelector subsetSelector: lbSubsetConfig.getSubsetSelectorsList()) { + final ProtocolStringList keys = subsetSelector.getKeysList(); + for (Endpoint endpoint : prioritySet.endpoints()) { + final LbEndpoint lbEndpoint = endpoint.attr(XdsAttributeKeys.LB_ENDPOINT_KEY); + assert lbEndpoint != null; + final Struct endpointMetadata = lbEndpoint.getMetadata().getFilterMetadataOrDefault( + SUBSET_LOAD_BALANCING_FILTER_NAME, Struct.getDefaultInstance()); + final Struct.Builder filteredStructBuilder = Struct.newBuilder(); + boolean allKeysFound = true; + for (String key : keys) { + if (!endpointMetadata.containsFields(key)) { + allKeysFound = false; + break; + } + filteredStructBuilder.putFields(key, endpointMetadata.getFieldsOrThrow(key)); + } + if (!allKeysFound) { + continue; + } + final Struct filteredStruct = filteredStructBuilder.build(); + endpointsPerFilterStruct.computeIfAbsent(filteredStruct, unused -> new ArrayList<>()) + .add(endpoint); + } + } + final ImmutableMap.Builder builder = ImmutableMap.builder(); + for (Entry> entry : endpointsPerFilterStruct.entrySet()) { + final PrioritySet subsetPrioritySet = + new PriorityStateManager(clusterSnapshot, entry.getValue()).build(); + final DefaultLoadBalancer subsetLoadBalancer = + new DefaultLoadBalancer(subsetPrioritySet, localCluster, localPrioritySet); + builder.put(entry.getKey(), subsetLoadBalancer); + } + return builder.build(); } @Override public String toString() { return MoreObjects.toStringHelper(this) - .add("loadBalancer", loadBalancer) + .add("lbSubsetConfig", lbSubsetConfig) + .add("subsetLoadBalancers", subsetLoadBalancers) + .add("allEndpointsLoadBalancer", allEndpointsLoadBalancer) .toString(); } @@ -116,10 +168,4 @@ public String toString() { public PrioritySet prioritySet() { return prioritySet; } - - @Override - @Nullable - public LocalityRoutingState localityRoutingState() { - return localityRoutingState; - } } diff --git a/xds/src/main/java/com/linecorp/armeria/xds/client/endpoint/XdsAttributeKeys.java b/xds/src/main/java/com/linecorp/armeria/xds/client/endpoint/XdsAttributeKeys.java index 83437321e97..a5658f3083d 100644 --- a/xds/src/main/java/com/linecorp/armeria/xds/client/endpoint/XdsAttributeKeys.java +++ b/xds/src/main/java/com/linecorp/armeria/xds/client/endpoint/XdsAttributeKeys.java @@ -16,6 +16,7 @@ package com.linecorp.armeria.xds.client.endpoint; +import io.envoyproxy.envoy.config.core.v3.Metadata; import io.envoyproxy.envoy.config.endpoint.v3.LbEndpoint; import io.envoyproxy.envoy.config.endpoint.v3.LocalityLbEndpoints; import io.netty.util.AttributeKey; @@ -28,6 +29,8 @@ final class XdsAttributeKeys { AttributeKey.valueOf(XdsAttributeKeys.class, "LOCALITY_LB_ENDPOINTS_KEY"); static final AttributeKey XDS_RANDOM = AttributeKey.valueOf(XdsAttributeKeys.class, "XDS_RANDOM"); + public static final AttributeKey ROUTE_METADATA_MATCH = + AttributeKey.valueOf(XdsAttributeKeys.class, "ROUTER_METADATA"); private XdsAttributeKeys() {} } diff --git a/xds/src/main/java/com/linecorp/armeria/xds/client/endpoint/XdsLoadBalancer.java b/xds/src/main/java/com/linecorp/armeria/xds/client/endpoint/XdsLoadBalancer.java index 3bfa8e07417..ffde26bc380 100644 --- a/xds/src/main/java/com/linecorp/armeria/xds/client/endpoint/XdsLoadBalancer.java +++ b/xds/src/main/java/com/linecorp/armeria/xds/client/endpoint/XdsLoadBalancer.java @@ -16,16 +16,7 @@ package com.linecorp.armeria.xds.client.endpoint; -import com.google.common.annotations.VisibleForTesting; - -import com.linecorp.armeria.common.annotation.Nullable; -import com.linecorp.armeria.xds.client.endpoint.LocalityRoutingStateFactory.LocalityRoutingState; - interface XdsLoadBalancer extends LoadBalancer { PrioritySet prioritySet(); - - @Nullable - @VisibleForTesting - LocalityRoutingState localityRoutingState(); } diff --git a/xds/src/test/java/com/linecorp/armeria/xds/AggregatingNodeTest.java b/xds/src/test/java/com/linecorp/armeria/xds/AggregatingNodeTest.java index c7d805c087d..bea693b6374 100644 --- a/xds/src/test/java/com/linecorp/armeria/xds/AggregatingNodeTest.java +++ b/xds/src/test/java/com/linecorp/armeria/xds/AggregatingNodeTest.java @@ -21,8 +21,6 @@ import java.net.URI; import java.util.Deque; -import java.util.List; -import java.util.Map.Entry; import java.util.concurrent.ConcurrentLinkedDeque; import java.util.concurrent.TimeUnit; @@ -40,6 +38,7 @@ import io.envoyproxy.controlplane.cache.v3.Snapshot; import io.envoyproxy.controlplane.server.V3DiscoveryServer; import io.envoyproxy.envoy.config.bootstrap.v3.Bootstrap; +import io.envoyproxy.envoy.config.cluster.v3.Cluster; import io.envoyproxy.envoy.config.cluster.v3.Cluster.DiscoveryType; import io.envoyproxy.envoy.config.endpoint.v3.ClusterLoadAssignment; import io.envoyproxy.envoy.config.listener.v3.Listener; @@ -230,20 +229,27 @@ private static void assertListenerSnapshot(ListenerSnapshot listenerSnapshot, St // verify route final RouteSnapshot routeSnapshot = listenerSnapshot.routeSnapshot(); final String routeName = routeSnapshot.xdsResource().resource().getName(); - final RouteConfiguration snapshotRouteConfiguration = snapshot.routes().resources().get(routeName); - assertThat(routeSnapshot.xdsResource().resource()).isEqualTo(snapshotRouteConfiguration); - - int i = 0; - for (Entry> e : routeSnapshot.virtualHostMap().entrySet()) { - final VirtualHost snapshotVirtualHost = snapshotRouteConfiguration.getVirtualHosts(i++); - assertThat(e.getKey()).isEqualTo(snapshotVirtualHost); - final List clusterSnapshots = e.getValue(); - for (int j = 0; j < clusterSnapshots.size(); j++) { - final ClusterSnapshot clusterSnapshot = clusterSnapshots.get(j); - assertThat(clusterSnapshot.route()).isEqualTo(snapshotVirtualHost.getRoutes(j)); - final String clusterName = clusterSnapshot.xdsResource().resource().getName(); - assertThat(clusterSnapshot.xdsResource().resource()) - .isEqualTo(snapshot.clusters().resources().get(clusterName)); + final RouteConfiguration expectedRoute = snapshot.routes().resources().get(routeName); + assertThat(routeSnapshot.xdsResource().resource()).isEqualTo(expectedRoute); + + for (int i = 0; i < routeSnapshot.virtualHostSnapshots().size(); i++) { + // validate virtual host + final VirtualHostSnapshot virtualHostSnapshot = routeSnapshot.virtualHostSnapshots().get(i); + final VirtualHost expectedVirtualHost = expectedRoute.getVirtualHosts(i); + assertThat(virtualHostSnapshot.xdsResource().resource()).isEqualTo(expectedVirtualHost); + for (int j = 0; j < virtualHostSnapshot.routeEntries().size(); j++) { + final RouteEntry routeEntry = virtualHostSnapshot.routeEntries().get(j); + assertThat(routeEntry.route()).isEqualTo(expectedVirtualHost.getRoutes(j)); + + // validate cluster + final ClusterSnapshot clusterSnapshot = routeEntry.clusterSnapshot(); + assertThat(clusterSnapshot).isNotNull(); + final String clusterName = clusterSnapshot.xdsResource().name(); + assertThat(clusterName).isEqualTo(expectedVirtualHost.getRoutes(j).getRoute().getCluster()); + final Cluster expectedCluster = snapshot.clusters().resources().get(clusterName); + assertThat(clusterSnapshot.xdsResource().resource()).isEqualTo(expectedCluster); + + // validate endpoint final EndpointSnapshot endpointSnapshot = clusterSnapshot.endpointSnapshot(); assertThat(endpointSnapshot.xdsResource().resource()) .isEqualTo(snapshot.endpoints().resources().get(clusterName)); diff --git a/xds/src/test/java/com/linecorp/armeria/xds/DynamicResourcesTest.java b/xds/src/test/java/com/linecorp/armeria/xds/DynamicResourcesTest.java index 95f87237c8b..03e7adf8d43 100644 --- a/xds/src/test/java/com/linecorp/armeria/xds/DynamicResourcesTest.java +++ b/xds/src/test/java/com/linecorp/armeria/xds/DynamicResourcesTest.java @@ -114,7 +114,9 @@ void basicCase() throws Exception { final Cluster expectedCluster = cache.getSnapshot(GROUP).clusters().resources().get(clusterName); - final ClusterSnapshot clusterSnapshot = routeSnapshot.clusterSnapshots().get(0); + final ClusterSnapshot clusterSnapshot = routeSnapshot.virtualHostSnapshots().get(0) + .routeEntries().get(0) + .clusterSnapshot(); assertThat(clusterSnapshot.xdsResource().resource()).isEqualTo(expectedCluster); final ClusterLoadAssignment expectedEndpoint = diff --git a/xds/src/test/java/com/linecorp/armeria/xds/MostlyStaticWithDynamicEdsTest.java b/xds/src/test/java/com/linecorp/armeria/xds/MostlyStaticWithDynamicEdsTest.java index 58adfa9ea49..1186ea59004 100644 --- a/xds/src/test/java/com/linecorp/armeria/xds/MostlyStaticWithDynamicEdsTest.java +++ b/xds/src/test/java/com/linecorp/armeria/xds/MostlyStaticWithDynamicEdsTest.java @@ -99,7 +99,9 @@ void basicCase() throws Exception { final RouteSnapshot routeSnapshot = listenerSnapshot.routeSnapshot(); assertThat(routeSnapshot.xdsResource().resource()).isEqualTo(expectedRoute); - final ClusterSnapshot clusterSnapshot = routeSnapshot.clusterSnapshots().get(0); + final ClusterSnapshot clusterSnapshot = routeSnapshot.virtualHostSnapshots().get(0) + .routeEntries().get(0) + .clusterSnapshot(); assertThat(clusterSnapshot.xdsResource().resource()).isEqualTo(staticCluster); final ClusterLoadAssignment expectedEndpoint = cache.getSnapshot(GROUP).endpoints().resources().get("cluster"); diff --git a/xds/src/test/java/com/linecorp/armeria/xds/MultiConfigSourceTest.java b/xds/src/test/java/com/linecorp/armeria/xds/MultiConfigSourceTest.java index 7c3e2cd2c40..a44030c3586 100644 --- a/xds/src/test/java/com/linecorp/armeria/xds/MultiConfigSourceTest.java +++ b/xds/src/test/java/com/linecorp/armeria/xds/MultiConfigSourceTest.java @@ -155,8 +155,11 @@ void fromListener() throws Exception { // Updates are propagated for the initial value final ClusterLoadAssignment expected = cache2.getSnapshot(GROUP).endpoints().resources().get("cluster1"); - assertThat(listenerSnapshot.routeSnapshot().clusterSnapshots() - .get(0).endpointSnapshot().xdsResource().resource()).isEqualTo(expected); + assertThat(listenerSnapshot.routeSnapshot() + .virtualHostSnapshots().get(0) + .routeEntries().get(0) + .clusterSnapshot().endpointSnapshot().xdsResource().resource()) + .isEqualTo(expected); await().pollDelay(100, TimeUnit.MILLISECONDS) .untilAsserted(() -> assertThat(watcher.events()).isEmpty()); @@ -175,8 +178,11 @@ void basicSelfConfigSource() { // Updates are propagated for the initial value final ClusterLoadAssignment expected = cache1.getSnapshot(GROUP).endpoints().resources().get("self-cluster1"); - assertThat(listenerSnapshot.routeSnapshot().clusterSnapshots() - .get(0).endpointSnapshot().xdsResource().resource()).isEqualTo(expected); + assertThat(listenerSnapshot.routeSnapshot() + .virtualHostSnapshots().get(0) + .routeEntries().get(0) + .clusterSnapshot().endpointSnapshot().xdsResource().resource()) + .isEqualTo(expected); await().pollDelay(100, TimeUnit.MILLISECONDS) .untilAsserted(() -> assertThat(watcher.events()).isEmpty()); @@ -195,8 +201,11 @@ void adsSelfConfigSource() { // Updates are propagated for the initial value final ClusterLoadAssignment expected = cache2.getSnapshot(GROUP).endpoints().resources().get("self-cluster2"); - assertThat(listenerSnapshot.routeSnapshot().clusterSnapshots() - .get(0).endpointSnapshot().xdsResource().resource()).isEqualTo(expected); + assertThat(listenerSnapshot.routeSnapshot() + .virtualHostSnapshots().get(0) + .routeEntries().get(0) + .clusterSnapshot().endpointSnapshot().xdsResource().resource()) + .isEqualTo(expected); await().pollDelay(100, TimeUnit.MILLISECONDS) .untilAsserted(() -> assertThat(watcher.events()).isEmpty()); diff --git a/xds/src/test/java/com/linecorp/armeria/xds/client/endpoint/RouteMetadataSubsetTest.java b/xds/src/test/java/com/linecorp/armeria/xds/client/endpoint/RouteMetadataSubsetTest.java index c99424ed7e2..7b98a0a8142 100644 --- a/xds/src/test/java/com/linecorp/armeria/xds/client/endpoint/RouteMetadataSubsetTest.java +++ b/xds/src/test/java/com/linecorp/armeria/xds/client/endpoint/RouteMetadataSubsetTest.java @@ -163,22 +163,6 @@ configSource, staticResourceListener(routeMetadataMatch1, clusterName), assertThat(endpointGroup.selectNow(ctx)).isEqualTo(Endpoint.of("127.0.0.1", 8082)); } - // No Route metadata so use all endpoints. - final Metadata routeMetadataMatch2 = Metadata.newBuilder().putFilterMetadata( - SUBSET_LOAD_BALANCING_FILTER_NAME, Struct.getDefaultInstance()).build(); - bootstrap = XdsTestResources.bootstrap(configSource, - staticResourceListener(routeMetadataMatch2, clusterName), - bootstrapCluster); - try (XdsBootstrap xdsBootstrap = XdsBootstrap.of(bootstrap); - EndpointGroup endpointGroup = XdsEndpointGroup.of("listener", xdsBootstrap)) { - - await().untilAsserted(() -> assertThat(endpointGroup.whenReady()).isDone()); - final ClientRequestContext ctx = ClientRequestContext.of(HttpRequest.of(HttpMethod.GET, "/")); - assertThat(endpointGroup.selectNow(ctx)).isEqualTo(Endpoint.of("127.0.0.1", 8080)); - assertThat(endpointGroup.selectNow(ctx)).isEqualTo(Endpoint.of("127.0.0.1", 8081)); - assertThat(endpointGroup.selectNow(ctx)).isEqualTo(Endpoint.of("127.0.0.1", 8082)); - } - final Metadata routeMetadataMatch3 = Metadata.newBuilder().putFilterMetadata( SUBSET_LOAD_BALANCING_FILTER_NAME, Struct.newBuilder() .putFields("foo", stringValue("foo1")) diff --git a/xds/src/test/java/com/linecorp/armeria/xds/client/endpoint/ZoneAwareTest.java b/xds/src/test/java/com/linecorp/armeria/xds/client/endpoint/ZoneAwareTest.java index 7abebe72261..703713f6f6e 100644 --- a/xds/src/test/java/com/linecorp/armeria/xds/client/endpoint/ZoneAwareTest.java +++ b/xds/src/test/java/com/linecorp/armeria/xds/client/endpoint/ZoneAwareTest.java @@ -100,7 +100,6 @@ void pickLocalRegion(int upstreamHealthyHosts, int upstreamAHealthyHosts, final ClusterEntry clusterEntry = endpointGroup.clusterEntriesMap().get("cluster"); assertThat(clusterEntry).isNotNull(); assertThat(clusterEntry.latestValue()).isNotNull(); - assertThat(clusterEntry.latestValue().localityRoutingState()).isNotNull(); }); final ClientRequestContext ctx = ClientRequestContext.of(HttpRequest.of(HttpMethod.GET, "/")); final SettableXdsRandom random = new SettableXdsRandom(); @@ -157,7 +156,6 @@ void routingEnabled(int routingEnabled, String expectedRegionsStr) throws Except final ClusterEntry clusterEntry = endpointGroup.clusterEntriesMap().get("cluster"); assertThat(clusterEntry).isNotNull(); assertThat(clusterEntry.latestValue()).isNotNull(); - assertThat(clusterEntry.latestValue().localityRoutingState()).isNotNull(); }); final ClientRequestContext ctx = ClientRequestContext.of(HttpRequest.of(HttpMethod.GET, "/")); final SettableXdsRandom random = new SettableXdsRandom(); @@ -214,7 +212,6 @@ void panicMode(int threshold, String expectedRegionsStr) throws Exception { final ClusterEntry clusterEntry = endpointGroup.clusterEntriesMap().get("cluster"); assertThat(clusterEntry).isNotNull(); assertThat(clusterEntry.latestValue()).isNotNull(); - assertThat(clusterEntry.latestValue().localityRoutingState()).isNotNull(); }); final ClientRequestContext ctx = ClientRequestContext.of(HttpRequest.of(HttpMethod.GET, "/")); final SettableXdsRandom random = new SettableXdsRandom(); @@ -268,7 +265,6 @@ void pickResidualRegion(long localPercentage, String expectedRegion) throws Exce final ClusterEntry clusterEntry = endpointGroup.clusterEntriesMap().get("cluster"); assertThat(clusterEntry).isNotNull(); assertThat(clusterEntry.latestValue()).isNotNull(); - assertThat(clusterEntry.latestValue().localityRoutingState()).isNotNull(); }); final ClientRequestContext ctx = ClientRequestContext.of(HttpRequest.of(HttpMethod.GET, "/")); final SettableXdsRandom random = new SettableXdsRandom(); @@ -324,7 +320,6 @@ void multiResidualRegion(long localThreshold, String expectedRegion) throws Exce final ClusterEntry clusterEntry = endpointGroup.clusterEntriesMap().get("cluster"); assertThat(clusterEntry).isNotNull(); assertThat(clusterEntry.latestValue()).isNotNull(); - assertThat(clusterEntry.latestValue().localityRoutingState()).isNotNull(); }); final ClientRequestContext ctx = ClientRequestContext.of(HttpRequest.of(HttpMethod.GET, "/")); final SettableXdsRandom random = new SettableXdsRandom(); @@ -384,7 +379,6 @@ void onlyUpstreamResidual(long localThreshold, String expectedRegion) throws Exc final ClusterEntry clusterEntry = endpointGroup.clusterEntriesMap().get("cluster"); assertThat(clusterEntry).isNotNull(); assertThat(clusterEntry.latestValue()).isNotNull(); - assertThat(clusterEntry.latestValue().localityRoutingState()).isNotNull(); }); final ClientRequestContext ctx = ClientRequestContext.of(HttpRequest.of(HttpMethod.GET, "/")); final SettableXdsRandom random = new SettableXdsRandom(); @@ -444,7 +438,6 @@ void priorityOneNotUsed(int priority, int selectPriority, String expectedRegion) final ClusterEntry clusterEntry = endpointGroup.clusterEntriesMap().get("cluster"); assertThat(clusterEntry).isNotNull(); assertThat(clusterEntry.latestValue()).isNotNull(); - assertThat(clusterEntry.latestValue().localityRoutingState()).isNotNull(); }); final ClientRequestContext ctx = ClientRequestContext.of(HttpRequest.of(HttpMethod.GET, "/")); final SettableXdsRandom random = new SettableXdsRandom(); @@ -509,7 +502,6 @@ void zoneAwareConfigurationsAreRespected(ZoneAwareLbConfig zoneAwareLbConfig, final ClusterEntry clusterEntry = endpointGroup.clusterEntriesMap().get("cluster"); assertThat(clusterEntry).isNotNull(); assertThat(clusterEntry.latestValue()).isNotNull(); - assertThat(clusterEntry.latestValue().localityRoutingState()).isNotNull(); }); final ClientRequestContext ctx = ClientRequestContext.of(HttpRequest.of(HttpMethod.GET, "/")); final SettableXdsRandom random = new SettableXdsRandom(); From 0d230732842e0239e3b9faa22ed39ab8f50d1799 Mon Sep 17 00:00:00 2001 From: jrhee17 Date: Wed, 26 Mar 2025 11:47:39 +0900 Subject: [PATCH 2/2] address comments by @minwoox and @trustin --- .../linecorp/armeria/xds/ClusterResourceNode.java | 5 ----- .../java/com/linecorp/armeria/xds/RouteEntry.java | 2 +- .../com/linecorp/armeria/xds/RouteResourceNode.java | 6 ++++-- .../armeria/xds/VirtualHostResourceNode.java | 1 + .../linecorp/armeria/xds/VirtualHostSnapshot.java | 5 +++-- .../linecorp/armeria/xds/VirtualHostXdsResource.java | 12 +++++++++++- .../xds/client/endpoint/XdsAttributeKeys.java | 2 +- 7 files changed, 21 insertions(+), 12 deletions(-) diff --git a/xds/src/main/java/com/linecorp/armeria/xds/ClusterResourceNode.java b/xds/src/main/java/com/linecorp/armeria/xds/ClusterResourceNode.java index fb8d69835b4..b95b1e1ad62 100644 --- a/xds/src/main/java/com/linecorp/armeria/xds/ClusterResourceNode.java +++ b/xds/src/main/java/com/linecorp/armeria/xds/ClusterResourceNode.java @@ -80,11 +80,6 @@ public void doOnChanged(ClusterXdsResource resource) { } } - @Override - public void close() { - super.close(); - } - private class EndpointSnapshotWatcher implements SnapshotWatcher { @Override public void snapshotUpdated(EndpointSnapshot newSnapshot) { diff --git a/xds/src/main/java/com/linecorp/armeria/xds/RouteEntry.java b/xds/src/main/java/com/linecorp/armeria/xds/RouteEntry.java index 892abc52704..077d8b96688 100644 --- a/xds/src/main/java/com/linecorp/armeria/xds/RouteEntry.java +++ b/xds/src/main/java/com/linecorp/armeria/xds/RouteEntry.java @@ -1,5 +1,5 @@ /* - * Copyright 2025 LINE Corporation + * Copyright 2025 LY 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 diff --git a/xds/src/main/java/com/linecorp/armeria/xds/RouteResourceNode.java b/xds/src/main/java/com/linecorp/armeria/xds/RouteResourceNode.java index 35f325095c2..8d6e2c78093 100644 --- a/xds/src/main/java/com/linecorp/armeria/xds/RouteResourceNode.java +++ b/xds/src/main/java/com/linecorp/armeria/xds/RouteResourceNode.java @@ -50,13 +50,15 @@ final class RouteResourceNode extends AbstractResourceNodeWithPrimer virtualHosts = routeConfiguration.getVirtualHostsList(); + for (int i = 0; i < virtualHosts.size(); i++) { pending.add(i); virtualHostSnapshots.add(null); - final VirtualHost virtualHost = routeConfiguration.getVirtualHostsList().get(i); + final VirtualHost virtualHost = virtualHosts.get(i); final VirtualHostResourceNode childNode = StaticResourceUtils.staticVirtualHost(xdsBootstrap(), name(), resource, snapshotWatcher, i, virtualHost); diff --git a/xds/src/main/java/com/linecorp/armeria/xds/VirtualHostResourceNode.java b/xds/src/main/java/com/linecorp/armeria/xds/VirtualHostResourceNode.java index 20c2e6c4c56..3a9322bc30b 100644 --- a/xds/src/main/java/com/linecorp/armeria/xds/VirtualHostResourceNode.java +++ b/xds/src/main/java/com/linecorp/armeria/xds/VirtualHostResourceNode.java @@ -56,6 +56,7 @@ final class VirtualHostResourceNode extends AbstractResourceNodeWithPrimer clusterSnapshots, int index) { this.virtualHostXdsResource = virtualHostXdsResource; - assert clusterSnapshots.size() == virtualHostXdsResource.resource().getRoutesCount(); + final VirtualHost virtualHost = virtualHostXdsResource.resource(); + assert clusterSnapshots.size() == virtualHost.getRoutesCount(); final ImmutableList.Builder routeEntriesBuilder = ImmutableList.builder(); for (int i = 0; i < clusterSnapshots.size(); i++) { final ClusterSnapshot clusterSnapshot = clusterSnapshots.get(i); - final Route route = virtualHostXdsResource.resource().getRoutes(i); + final Route route = virtualHost.getRoutes(i); routeEntriesBuilder.add(new RouteEntry(route, clusterSnapshot)); } routeEntries = routeEntriesBuilder.build(); diff --git a/xds/src/main/java/com/linecorp/armeria/xds/VirtualHostXdsResource.java b/xds/src/main/java/com/linecorp/armeria/xds/VirtualHostXdsResource.java index adeaf701e82..aa6a4ba65ea 100644 --- a/xds/src/main/java/com/linecorp/armeria/xds/VirtualHostXdsResource.java +++ b/xds/src/main/java/com/linecorp/armeria/xds/VirtualHostXdsResource.java @@ -16,6 +16,8 @@ package com.linecorp.armeria.xds; +import com.google.common.base.MoreObjects; + import com.linecorp.armeria.common.annotation.Nullable; import io.envoyproxy.envoy.config.route.v3.VirtualHost; @@ -34,7 +36,7 @@ public final class VirtualHostXdsResource extends XdsResourceWithPrimer XDS_RANDOM = AttributeKey.valueOf(XdsAttributeKeys.class, "XDS_RANDOM"); public static final AttributeKey ROUTE_METADATA_MATCH = - AttributeKey.valueOf(XdsAttributeKeys.class, "ROUTER_METADATA"); + AttributeKey.valueOf(XdsAttributeKeys.class, "ROUTER_METADATA_MATCH"); private XdsAttributeKeys() {} }