diff --git a/packages/grpc-js-xds/src/xds-client.ts b/packages/grpc-js-xds/src/xds-client.ts index 785c11943..3c496e406 100644 --- a/packages/grpc-js-xds/src/xds-client.ts +++ b/packages/grpc-js-xds/src/xds-client.ts @@ -150,6 +150,9 @@ class ResourceTimer { if (!resourceState) { return; } + if (resourceState.cachedResource !== null) { + return; + } resourceState.meta.clientStatus = 'DOES_NOT_EXIST'; for (const watcher of resourceState.watchers) { watcher.onResourceDoesNotExist(); @@ -372,6 +375,7 @@ class AdsCallState { experimental.log(logVerbosity.ERROR, 'Ignoring nonexistent resource ' + xdsResourceNameToString({authority, key}, result.type!.getTypeUrl())); resourceState.deletionIgnored = true; } else { + resourceState.cachedResource = null; resourceState.meta.clientStatus = 'DOES_NOT_EXIST'; process.nextTick(() => { for (const watcher of resourceState.watchers) { @@ -404,6 +408,13 @@ class AdsCallState { this.trace( 'ADS stream ended. code=' + streamStatus.code + ' details= ' + streamStatus.details ); + for (const typeState of this.typeStates.values()) { + for (const authorityMap of typeState.subscribedResources.values()) { + for (const timer of authorityMap.values()) { + timer.maybeCancelTimer(); + } + } + } if (streamStatus.code !== status.OK && !this.receivedAnyResponse) { for (const watcher of this.allWatchers()) { watcher.onError(streamStatus); @@ -458,6 +469,7 @@ class AdsCallState { if (!authorityMap) { return; } + authorityMap.get(name.key)?.maybeCancelTimer(); authorityMap.delete(name.key); if (authorityMap.size === 0) { typeState.subscribedResources.delete(name.authority); @@ -937,6 +949,9 @@ class XdsSingleServerClient { const metadata = new Metadata({waitForReady: true}); const call = this.adsClient.StreamAggregatedResources(metadata); this.adsCallState = new AdsCallState(this, call, this.xdsClient.adsNode!); + if (this.adsClient.getChannel().getConnectivityState(false) === connectivityState.READY) { + this.adsCallState.markStreamStarted(); + } this.adsBackoff.runOnce(); } diff --git a/packages/grpc-js-xds/test/test-core.ts b/packages/grpc-js-xds/test/test-core.ts index 3eef81dd9..e6b369b48 100644 --- a/packages/grpc-js-xds/test/test-core.ts +++ b/packages/grpc-js-xds/test/test-core.ts @@ -197,5 +197,46 @@ describe('core xDS functionality', () => { xdsServer.setRdsResource(routeGroup2.getRouteConfiguration()); await cluster2.waitForAllBackendsToReceiveTraffic(); client.stopCalls(); - }) + }); + it('should recover when a deleted LDS resource is restored with identical content', async () => { + const [backend] = await createBackends(1); + const serverRoute = new FakeServerRoute(backend.getPort(), 'serverRoute'); + xdsServer.setRdsResource(serverRoute.getRouteConfiguration()); + xdsServer.setLdsResource(serverRoute.getListener()); + xdsServer.addResponseListener((typeUrl, responseState) => { + if (responseState.state === 'NACKED') { + client?.stopCalls(); + assert.fail(`Client NACKED ${typeUrl} resource with message ${responseState.errorMessage}`); + } + }); + const cluster = new FakeEdsCluster('cluster1', 'endpoint1', [{backends: [backend], locality: {region: 'region1'}}]); + const routeGroup = new FakeRouteGroup('listener1', 'route1', [{cluster: cluster}]); + await routeGroup.startAllBackends(xdsServer); + xdsServer.setEdsResource(cluster.getEndpointConfig()); + xdsServer.setCdsResource(cluster.getClusterConfig()); + xdsServer.setRdsResource(routeGroup.getRouteConfiguration()); + xdsServer.setLdsResource(routeGroup.getListener()); + client = XdsTestClient.createFromServer('listener1', xdsServer); + client.startCalls(100); + await routeGroup.waitForAllBackendsToReceiveTraffic(); + client.stopCalls(); + + xdsServer.unsetLdsResource('listener1'); + + const deadline = Date.now() + 1000; + while (client.getConnectivityState() === connectivityState.READY && Date.now() < deadline) { + await new Promise(resolve => setTimeout(resolve, 50)); + } + + const error = await client.sendOneCallAsync(); + assert(error, 'Expected RPC to fail after LDS deletion'); + + // Restore identical LDS resource + xdsServer.setLdsResource(routeGroup.getListener()); + + // Verify client recovers and traffic flows normally + client.startCalls(100); + await routeGroup.waitForAllBackendsToReceiveTraffic(); + client.stopCalls(); + }); }); diff --git a/packages/grpc-js-xds/test/xds-server.ts b/packages/grpc-js-xds/test/xds-server.ts index d6115e79a..e3eb48066 100644 --- a/packages/grpc-js-xds/test/xds-server.ts +++ b/packages/grpc-js-xds/test/xds-server.ts @@ -177,6 +177,11 @@ export class ControlPlaneServer { this.setResource({...resource, '@type': LDS_TYPE_URL}, resource.name!); } + unsetLdsResource(name: string) { + trace(`unsetLdsResource(${name})`); + this.unsetResource(LDS_TYPE_URL, name); + } + setRdsResource(resource: RouteConfiguration) { trace(`setRdsResource(${resource.name!})`); this.setResource({...resource, '@type': RDS_TYPE_URL}, resource.name!); @@ -209,7 +214,7 @@ export class ControlPlaneServer { private sendResourceUpdates(typeUrl: T, clients: Set, includeResources: Set) { const resourceTypeState = this.resourceMap[typeUrl] as ResourceTypeState; - const clientResources = new Map(); + const clientResources = new Map(Array.from(clients, client => [client, []])); for (const [resourceName, resourceState] of resourceTypeState.resourceNameMap) { /* For RDS and EDS, only send updates for the listed updated resources. * Otherwise include all resources. */