The latest gRPC library does unsubscribe resource if channel is not used and then re-subscribed. This not work properly when using V3DiscoveryServer with SimpleCache<>.
The question is how it should look like protocol flow in the SotW protocol variant.
Here is whats going on:
- subscribe resource "a" and "b"
- unsubscribe resource "b"
- subscribe resource "a" and "b": Here library didn't response, however the xds-client already removed "b" and expecting response. If instead of re-subscribing "b" is subscribed "c", the control-plain responds.

According https://www.envoyproxy.io/docs/envoy/latest/api-docs/xds_protocol#unsubscribing-from-resources
Unsubscribing From Resources
In the SotW protocol variants, each request must contain the full list of resource names being subscribed to in the resource_names field, so unsubscribing to a set of resources is done by sending a new request containing all resource names that are still being subscribed to but not containing the resource names being unsubscribed to. For example, if the client had previously been subscribed to resources A and B but wishes to unsubscribe from B, it must send a new request containing only resource A.
Looks like that gRPC confirms the specification at un-subscribe.
Acording to specification https://www.envoyproxy.io/docs/envoy/latest/api-docs/xds_protocol#how-the-client-specifies-what-resources-to-return
When the client sends a new request that changes the set of resources being requested, the server must resend any newly requested resources, even if it previously sent those resources without having been asked for them and the resources have not changed since that time.
Here i cannot tell. The resource "b" is not exactly new. It is already known resource but unsubscribed.
This unit test I used to see how it is implemented.
packagecz.seznam.profile.xds;
importcom.google.common.collect.ImmutableList;
importcom.google.protobuf.InvalidProtocolBufferException;
importio.envoyproxy.controlplane.cache.NodeGroup;
importio.envoyproxy.controlplane.cache.v3.SimpleCache;
importio.envoyproxy.controlplane.cache.v3.Snapshot;
importio.envoyproxy.controlplane.server.V3DiscoveryServer;
importio.envoyproxy.envoy.config.cluster.v3.Cluster;
importio.envoyproxy.envoy.config.core.v3.Node;
importio.envoyproxy.envoy.config.endpoint.v3.ClusterLoadAssignment;
importio.envoyproxy.envoy.config.listener.v3.Listener;
importio.envoyproxy.envoy.config.route.v3.RouteConfiguration;
importio.envoyproxy.envoy.extensions.transport_sockets.tls.v3.Secret;
importio.envoyproxy.envoy.service.discovery.v3.AggregatedDiscoveryServiceGrpc;
importio.envoyproxy.envoy.service.discovery.v3.DiscoveryRequest;
importio.envoyproxy.envoy.service.discovery.v3.DiscoveryResponse;
importio.envoyproxy.controlplane.cache.Resources;
importio.grpc.inprocess.InProcessChannelBuilder;
importio.grpc.inprocess.InProcessServerBuilder;
importio.grpc.stub.StreamObserver;
importjava.io.IOException;
importjava.util.ArrayList;
importjava.util.List;
importorg.junit.jupiter.api.Assertions;
importorg.junit.jupiter.api.BeforeAll;
importorg.junit.jupiter.api.Test;
publicclassReSubscribeTest {
static {
varrootLogger = org.apache.log4j.Logger.getRootLogger();
rootLogger.setLevel(org.apache.log4j.Level.DEBUG);
rootLogger.addAppender(neworg.apache.log4j.ConsoleAppender(
neworg.apache.log4j.PatternLayout("%-6r [%p] %c - %m%n")));
}
publicstaticclassStaticNodeGroupimplementsNodeGroup<String> {
@OverridepublicStringhash(Nodenode) {
returnnode.getCluster();
}
}
@BeforeAllpublicstaticvoidsetup() throwsException {
}
privatestaticfinalStringCLUSTER_NAME = "default-cluster";
staticfinalio.envoyproxy.envoy.config.core.v3.Nodenode = io.envoyproxy.envoy.config.core.v3.Node.newBuilder()
.setCluster("default-cluster").setUserAgentName("gRPC C-core linux")
.setUserAgentVersion("C-core 37.0.0").addClientFeatures("envoy.lb.does_not_support_overprovisioning")
.addClientFeatures("xds.config.resource-in-sotw").build();
staticfinalDiscoveryRequestdrListener1 = DiscoveryRequest.newBuilder()
.setTypeUrl(Resources.V3.LISTENER_TYPE_URL).setNode(node)
.addResourceNames("my-echo-server")
.addResourceNames("my-echo-server2")
.build();
staticfinalDiscoveryRequestdrRoute1 = DiscoveryRequest.newBuilder()
.setTypeUrl(Resources.V3.ROUTE_TYPE_URL).setNode(node)
.addResourceNames("my-echo-server-route")
.addResourceNames("my-echo-server2-route")
.build();
staticfinalDiscoveryRequestdrCluster1 = DiscoveryRequest.newBuilder()
.setTypeUrl(Resources.V3.CLUSTER_TYPE_URL).setNode(node)
.addResourceNames("my-echo-server-cluster")
.addResourceNames("my-echo-server2-cluster")
.build();
staticfinalDiscoveryRequestdrEndpoint1 = DiscoveryRequest.newBuilder()
.setTypeUrl(Resources.V3.ENDPOINT_TYPE_URL).setNode(node)
.addResourceNames("my-echo-server-cluster")
.addResourceNames("my-echo-server2-cluster")
.build();
privatestaticfinalStringVERSION1 = "88eb5ecd-1e5f-484f-9522-d5be08dd166a";
@SuppressWarnings("null")
privatestaticfinalSnapshotSNAPSHOT1 = Snapshot.create(
ImmutableList.of(Cluster.newBuilder().setName("my-echo-server-cluster").build(),
Cluster.newBuilder().setName("my-echo-server2-cluster").build(),
Cluster.newBuilder().setName("my-echo-server3-cluster").build()),
ImmutableList.of(ClusterLoadAssignment.newBuilder().setClusterName("my-echo-server-cluster").build(),
ClusterLoadAssignment.newBuilder().setClusterName("my-echo-server2-cluster").build(),
ClusterLoadAssignment.newBuilder().setClusterName("my-echo-server3-cluster").build()),
ImmutableList.of(Listener.newBuilder().setName("my-echo-server").build(),
Listener.newBuilder().setName("my-echo-server2").build(),
Listener.newBuilder().setName("my-echo-server3").build()),
ImmutableList.of(RouteConfiguration.newBuilder().setName("my-echo-server-route").build(),
RouteConfiguration.newBuilder().setName("my-echo-server2-route").build(),
RouteConfiguration.newBuilder().setName("my-echo-server3-route").build()),
ImmutableList.<Secret>of(),
VERSION1);
@Testvoidtest() throwsIOException {
varcache = newSimpleCache<>(newStaticNodeGroup());
cache.setSnapshot("default-cluster", SNAPSHOT1);
vardiscoveryServer = newV3DiscoveryServer(cache);
finalvarserverName = InProcessServerBuilder.generateName();
finalvarserver = InProcessServerBuilder
.forName(serverName).directExecutor().addService(discoveryServer.getAggregatedDiscoveryServiceImpl())
.build();
server.start();
finalvarstub = AggregatedDiscoveryServiceGrpc
.newStub(InProcessChannelBuilder.forName(serverName).directExecutor().build());
varxdsConfig = newXdsConfig();
finalvarrequestObserver = stub.streamAggregatedResources(newResponseStreamObserver(xdsConfig));
// LDS->RDS->CDS->EDS// LDSrequestObserver.onNext(drListener1);
Assertions.assertEquals("0", xdsConfig.ldsNonce);
Assertions.assertEquals(List.of("my-echo-server", "my-echo-server2"), xdsConfig.ldsResourcesNames);
// RDSrequestObserver.onNext(drRoute1);
Assertions.assertEquals("1", xdsConfig.rdsNonce);
Assertions.assertEquals(List.of("my-echo-server-route", "my-echo-server2-route"), xdsConfig.rdsResourcesNames);
// CDSrequestObserver.onNext(drCluster1);
Assertions.assertEquals("2", xdsConfig.cdsNonce);
Assertions.assertEquals(List.of("my-echo-server-cluster", "my-echo-server2-cluster"),
xdsConfig.cdsResourcesNames);
// EDSrequestObserver.onNext(drEndpoint1);
Assertions.assertEquals("3", xdsConfig.edsNonce);
Assertions.assertEquals(List.of("my-echo-server-cluster", "my-echo-server2-cluster"),
xdsConfig.edsResourcesNames);
// ACK configurationfinalvardrListener2 = DiscoveryRequest.newBuilder(drListener1).setResponseNonce(xdsConfig.ldsNonce)
.setVersionInfo(xdsConfig.ldsVersionInfo).build();
finalvardrRoute2 = DiscoveryRequest.newBuilder(drRoute1).setResponseNonce(xdsConfig.rdsNonce)
.setVersionInfo(xdsConfig.rdsVersionInfo).build();
finalvardrCluster2 = DiscoveryRequest.newBuilder(drCluster1).setResponseNonce(xdsConfig.cdsNonce)
.setVersionInfo(xdsConfig.cdsVersionInfo).build();
finalvardrEndpoint2 = DiscoveryRequest.newBuilder(drEndpoint1).setResponseNonce(xdsConfig.edsNonce)
.setVersionInfo(xdsConfig.edsVersionInfo).build();
requestObserver.onNext(drListener2);
requestObserver.onNext(drRoute2);
requestObserver.onNext(drCluster2);
requestObserver.onNext(drEndpoint2);
// verify that no other change was receivedAssertions.assertEquals("0", xdsConfig.ldsNonce);
Assertions.assertEquals("1", xdsConfig.rdsNonce);
Assertions.assertEquals("2", xdsConfig.cdsNonce);
Assertions.assertEquals("3", xdsConfig.edsNonce);
// unsubscribe one resourceSystem.out.println("----------------------------------------------------------------------");
System.out.println(" unsubscribe one resource ");
System.out.println("----------------------------------------------------------------------");
finalvardrListener3 = DiscoveryRequest.newBuilder(drListener2).clearResourceNames()
.addResourceNames("my-echo-server").build();
finalvardrRoute3 = DiscoveryRequest.newBuilder(drRoute2).clearResourceNames()
.addResourceNames("my-echo-server-route").build();
finalvardrCluster3 = DiscoveryRequest.newBuilder(drCluster2).clearResourceNames()
.addResourceNames("my-echo-server-cluster").build();
finalvardrEndpoint3 = DiscoveryRequest.newBuilder(drEndpoint2).clearResourceNames()
.addResourceNames("my-echo-server-cluster").build();
requestObserver.onNext(drListener3);
requestObserver.onNext(drRoute3);
requestObserver.onNext(drCluster3);
requestObserver.onNext(drEndpoint3);
// verify replyAssertions.assertEquals("0", xdsConfig.ldsNonce);
Assertions.assertEquals("1", xdsConfig.rdsNonce);
Assertions.assertEquals("2", xdsConfig.cdsNonce);
Assertions.assertEquals("3", xdsConfig.edsNonce);
// no change was sendSystem.out.println("----------------------------------------------------------------------");
System.out.println(" re-subscribe the resource my-echo-server2");
System.out.println("----------------------------------------------------------------------");
finalvardrListener4 = DiscoveryRequest.newBuilder(drListener3)
.addResourceNames("my-echo-server2").build();
requestObserver.onNext(drListener4);
// verify replyAssertions.assertEquals("0", xdsConfig.ldsNonce);
// no change was send!System.out.println("----------------------------------------------------------------------");
System.out.println(" subscribe new resource my-echo-server3");
System.out.println("----------------------------------------------------------------------");
finalvardrListener5 = DiscoveryRequest.newBuilder(drListener3)
.addResourceNames("my-echo-server3").build();
requestObserver.onNext(drListener5);
// verify replyAssertions.assertEquals("4", xdsConfig.ldsNonce);
Assertions.assertEquals(List.of("my-echo-server", "my-echo-server3"),
xdsConfig.ldsResourcesNames);
// only "my-echo-server3" was sendrequestObserver.onCompleted();
server.shutdown();
try {
server.awaitTermination();
} catch (InterruptedExceptione) {
}
System.out.println("-- end --");
}
staticclassXdsConfig {
// The flow used in gRPC is `LDS->RDS->CDS->EDS`StringldsVersionInfo;
StringldsNonce;
List<String> ldsResourcesNames = List.of();
StringrdsVersionInfo;
StringrdsNonce;
List<String> rdsResourcesNames = List.of();
StringcdsVersionInfo;
StringcdsNonce;
List<String> cdsResourcesNames = List.of();
StringedsVersionInfo;
StringedsNonce;
List<String> edsResourcesNames = List.of();
}
staticclassResponseStreamObserverimplementsStreamObserver<DiscoveryResponse> {
XdsConfigxdsConfig;
publicResponseStreamObserver(XdsConfigxdsConfig) {
this.xdsConfig = xdsConfig;
}
@OverridepublicvoidonNext(DiscoveryResponsevalue) {
System.out.println(value.toString());
if (value.getTypeUrl().equals(Resources.V3.LISTENER_TYPE_URL)) {
xdsConfig.ldsVersionInfo = value.getVersionInfo();
xdsConfig.ldsNonce = value.getNonce();
varresourceNames = newArrayList<String>();
for (varany : value.getResourcesList()) {
try {
varlistener = any.unpack(Listener.class);
resourceNames.add(listener.getName());
} catch (InvalidProtocolBufferExceptione) {
e.printStackTrace();
}
}
xdsConfig.ldsResourcesNames = resourceNames;
} elseif (value.getTypeUrl().equals(Resources.V3.ROUTE_TYPE_URL)) {
xdsConfig.rdsVersionInfo = value.getVersionInfo();
xdsConfig.rdsNonce = value.getNonce();
varresourceNames = newArrayList<String>();
for (varany : value.getResourcesList()) {
try {
varroute = any.unpack(RouteConfiguration.class);
resourceNames.add(route.getName());
} catch (InvalidProtocolBufferExceptione) {
e.printStackTrace();
}
}
xdsConfig.rdsResourcesNames = resourceNames;
} elseif (value.getTypeUrl().equals(Resources.V3.CLUSTER_TYPE_URL)) {
xdsConfig.cdsVersionInfo = value.getVersionInfo();
xdsConfig.cdsNonce = value.getNonce();
varresourceNames = newArrayList<String>();
for (varany : value.getResourcesList()) {
try {
varcluster = any.unpack(Cluster.class);
resourceNames.add(cluster.getName());
} catch (InvalidProtocolBufferExceptione) {
e.printStackTrace();
}
}
xdsConfig.cdsResourcesNames = resourceNames;
} elseif (value.getTypeUrl().equals(Resources.V3.ENDPOINT_TYPE_URL)) {
xdsConfig.edsVersionInfo = value.getVersionInfo();
xdsConfig.edsNonce = value.getNonce();
varresourceNames = newArrayList<String>();
for (varany : value.getResourcesList()) {
try {
varcla = any.unpack(ClusterLoadAssignment.class);
resourceNames.add(cla.getClusterName());
} catch (InvalidProtocolBufferExceptione) {
e.printStackTrace();
}
}
xdsConfig.edsResourcesNames = resourceNames;
}
}
@OverridepublicvoidonError(Throwablet) {
}
@OverridepublicvoidonCompleted() {
}
}
}debug.log.txt
The latest gRPC library does unsubscribe resource if channel is not used and then re-subscribed. This not work properly when using
V3DiscoveryServerwithSimpleCache<>.The question is how it should look like protocol flow in the SotW protocol variant.
Here is whats going on:
According https://www.envoyproxy.io/docs/envoy/latest/api-docs/xds_protocol#unsubscribing-from-resources
Looks like that gRPC confirms the specification at un-subscribe.
Acording to specification https://www.envoyproxy.io/docs/envoy/latest/api-docs/xds_protocol#how-the-client-specifies-what-resources-to-return
Here i cannot tell. The resource "b" is not exactly new. It is already known resource but unsubscribed.
This unit test I used to see how it is implemented.
debug.log.txt