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);
         }
     }
 

Reply via email to