This is an automated email from the ASF dual-hosted git repository.
albumenj pushed a commit to branch 3.1
in repository https://gitbox.apache.org/repos/asf/dubbo.git
The following commit(s) were added to refs/heads/3.1 by this push:
new d84d098f3e fix dubbo proxyless xds connect, fix Connection closed
after GOAWAY, add xds observer retry create (#10544)
d84d098f3e is described below
commit d84d098f3e0f9a234ba913f516e0807d11bcdfbd
Author: wucheng1997 <[email protected]>
AuthorDate: Mon Oct 10 19:13:53 2022 +0800
fix dubbo proxyless xds connect, fix Connection closed after GOAWAY, add
xds observer retry create (#10544)
---
.../dubbo/registry/xds/util/protocol/AbstractProtocol.java | 13 ++++++++++++-
1 file changed, 12 insertions(+), 1 deletion(-)
diff --git
a/dubbo-xds/src/main/java/org/apache/dubbo/registry/xds/util/protocol/AbstractProtocol.java
b/dubbo-xds/src/main/java/org/apache/dubbo/registry/xds/util/protocol/AbstractProtocol.java
index 9b6687f54b..eb21cf3548 100644
---
a/dubbo-xds/src/main/java/org/apache/dubbo/registry/xds/util/protocol/AbstractProtocol.java
+++
b/dubbo-xds/src/main/java/org/apache/dubbo/registry/xds/util/protocol/AbstractProtocol.java
@@ -154,6 +154,11 @@ public abstract class AbstractProtocol<T, S extends
DeltaResource<T>> implements
// observer reused
StreamObserver<DiscoveryRequest> observer =
requestObserverMap.get(request);
+ if (observer == null) {
+ observer = xdsChannel.createDeltaDiscoveryRequest(new
ResponseObserver(request));
+ requestObserverMap.put(request, observer);
+ }
+
// send request to control panel
observer.onNext(buildDiscoveryRequest(names));
@@ -231,11 +236,17 @@ public abstract class AbstractProtocol<T, S extends
DeltaResource<T>> implements
@Override
public void onError(Throwable t) {
logger.error("xDS Client received error message! detail:", t);
+ clear();
}
@Override
public void onCompleted() {
- // ignore
+ logger.info("xDS Client completed, requestId: " + requestId);
+ clear();
+ }
+
+ private void clear() {
+ requestObserverMap.remove(requestId);
}
}