From cbcbac9ce606ce5822791d0afb1ca95c7f29c6d1 Mon Sep 17 00:00:00 2001 From: Albumen Kevin Date: Fri, 20 Jan 2023 14:01:14 +0800 Subject: [PATCH] Introduce Invocation#getInvokedInvokers to get invokers for ClusterFilter (#11359) --- .../support/AbstractClusterInvoker.java | 11 ++--- .../cluster/directory/MockDirInvocation.java | 19 +++++++-- .../support/AbstractClusterInvokerTest.java | 34 +++++++++++---- .../com/alibaba/dubbo/rpc/Invocation.java | 15 ++++++- .../com/alibaba/dubbo/rpc/RpcInvocation.java | 17 ++++++-- .../org/apache/dubbo/cache/CacheTest.java | 19 +++++++-- .../apache/dubbo/filter/LegacyInvocation.java | 16 +++++-- .../apache/dubbo/service/MockInvocation.java | 19 +++++++-- .../java/org/apache/dubbo/rpc/Invocation.java | 25 +++++++++-- .../org/apache/dubbo/rpc/RpcInvocation.java | 42 ++++++++++++------- .../apache/dubbo/rpc/RpcInvocationTest.java | 22 +++++++++- 11 files changed, 186 insertions(+), 53 deletions(-) diff --git a/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/support/AbstractClusterInvoker.java b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/support/AbstractClusterInvoker.java index c4435a44c6..3fdd1fbfff 100644 --- a/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/support/AbstractClusterInvoker.java +++ b/dubbo-cluster/src/main/java/org/apache/dubbo/rpc/cluster/support/AbstractClusterInvoker.java @@ -16,6 +16,11 @@ */ package org.apache.dubbo.rpc.cluster.support; +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.ThreadLocalRandom; +import java.util.concurrent.atomic.AtomicBoolean; + import org.apache.dubbo.common.URL; import org.apache.dubbo.common.Version; import org.apache.dubbo.common.config.Configuration; @@ -39,11 +44,6 @@ import org.apache.dubbo.rpc.model.ApplicationModel; import org.apache.dubbo.rpc.model.ScopeModelUtil; import org.apache.dubbo.rpc.support.RpcUtils; -import java.util.ArrayList; -import java.util.List; -import java.util.concurrent.ThreadLocalRandom; -import java.util.concurrent.atomic.AtomicBoolean; - import static org.apache.dubbo.common.constants.CommonConstants.DEFAULT_LOADBALANCE; import static org.apache.dubbo.common.constants.CommonConstants.DEFAULT_RESELECT_COUNT; import static org.apache.dubbo.common.constants.CommonConstants.ENABLE_CONNECTIVITY_VALIDATION; @@ -376,6 +376,7 @@ public abstract class AbstractClusterInvoker implements ClusterInvoker { if (ProfilerSwitch.isEnableSimpleProfiler()) { InvocationProfilerUtils.enterProfiler(invocation, "Invoker invoke. Target Address: " + invoker.getUrl().getAddress()); } + invocation.addInvokedInvoker(invoker); result = invoker.invoke(invocation); } finally { clearContext(invoker); diff --git a/dubbo-cluster/src/test/java/org/apache/dubbo/rpc/cluster/directory/MockDirInvocation.java b/dubbo-cluster/src/test/java/org/apache/dubbo/rpc/cluster/directory/MockDirInvocation.java index 3b40d9c818..66585bed0d 100644 --- a/dubbo-cluster/src/test/java/org/apache/dubbo/rpc/cluster/directory/MockDirInvocation.java +++ b/dubbo-cluster/src/test/java/org/apache/dubbo/rpc/cluster/directory/MockDirInvocation.java @@ -16,15 +16,16 @@ */ package org.apache.dubbo.rpc.cluster.directory; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.function.Consumer; + import org.apache.dubbo.rpc.AttachmentsAdapter; import org.apache.dubbo.rpc.Invocation; import org.apache.dubbo.rpc.Invoker; import org.apache.dubbo.rpc.model.ServiceModel; -import java.util.HashMap; -import java.util.Map; -import java.util.function.Consumer; - import static org.apache.dubbo.common.constants.CommonConstants.DUBBO_VERSION_KEY; import static org.apache.dubbo.common.constants.CommonConstants.GROUP_KEY; import static org.apache.dubbo.common.constants.CommonConstants.PATH_KEY; @@ -169,6 +170,16 @@ class MockDirInvocation implements Invocation { return (String) getObjectAttachment(key, defaultValue); } + @Override + public void addInvokedInvoker(Invoker invoker) { + + } + + @Override + public List> getInvokedInvokers() { + return null; + } + @Override public Object getObjectAttachment(String key, Object defaultValue) { Object result = attachments.get(key); diff --git a/dubbo-cluster/src/test/java/org/apache/dubbo/rpc/cluster/support/AbstractClusterInvokerTest.java b/dubbo-cluster/src/test/java/org/apache/dubbo/rpc/cluster/support/AbstractClusterInvokerTest.java index 74966bb338..665a9ef65d 100644 --- a/dubbo-cluster/src/test/java/org/apache/dubbo/rpc/cluster/support/AbstractClusterInvokerTest.java +++ b/dubbo-cluster/src/test/java/org/apache/dubbo/rpc/cluster/support/AbstractClusterInvokerTest.java @@ -16,6 +16,14 @@ */ package org.apache.dubbo.rpc.cluster.support; +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.atomic.AtomicLong; + import org.apache.dubbo.common.URL; import org.apache.dubbo.common.extension.ExtensionLoader; import org.apache.dubbo.common.utils.NetUtils; @@ -33,7 +41,6 @@ import org.apache.dubbo.rpc.cluster.filter.DemoService; import org.apache.dubbo.rpc.cluster.loadbalance.LeastActiveLoadBalance; import org.apache.dubbo.rpc.cluster.loadbalance.RandomLoadBalance; import org.apache.dubbo.rpc.cluster.loadbalance.RoundRobinLoadBalance; - import org.junit.jupiter.api.AfterAll; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.Assertions; @@ -43,13 +50,6 @@ import org.junit.jupiter.api.Disabled; import org.junit.jupiter.api.Test; import org.mockito.Mockito; -import java.util.ArrayList; -import java.util.HashMap; -import java.util.List; -import java.util.Map; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.atomic.AtomicLong; - import static org.apache.dubbo.common.constants.CommonConstants.CONSUMER; import static org.apache.dubbo.common.constants.CommonConstants.ENABLE_CONNECTIVITY_VALIDATION; import static org.apache.dubbo.common.constants.CommonConstants.INTERFACE_KEY; @@ -206,6 +206,24 @@ class AbstractClusterInvokerTest { } } + @Test + void testSelectedInvokers() { + cluster = new AbstractClusterInvoker(dic) { + @Override + protected Result doInvoke(Invocation invocation, List invokers, LoadBalance loadbalance) + throws RpcException { + checkInvokers(invokers, invocation); + Invoker invoker = select(loadbalance, invocation, invokers, null); + return invokeWithContext(invoker, invocation); + } + }; + + // invoke + cluster.invoke(invocation); + + Assertions.assertEquals(Collections.singletonList(invoker1), invocation.getInvokedInvokers()); + } + @Test void testSelect_Invokersize1() throws Exception { invokers.clear(); diff --git a/dubbo-compatible/src/main/java/com/alibaba/dubbo/rpc/Invocation.java b/dubbo-compatible/src/main/java/com/alibaba/dubbo/rpc/Invocation.java index c016e3b0ea..2eda40f495 100644 --- a/dubbo-compatible/src/main/java/com/alibaba/dubbo/rpc/Invocation.java +++ b/dubbo-compatible/src/main/java/com/alibaba/dubbo/rpc/Invocation.java @@ -17,13 +17,14 @@ package com.alibaba.dubbo.rpc; -import org.apache.dubbo.rpc.model.ServiceModel; - import java.util.Collections; import java.util.HashMap; +import java.util.List; import java.util.Map; import java.util.function.Consumer; +import org.apache.dubbo.rpc.model.ServiceModel; + @Deprecated public interface Invocation extends org.apache.dubbo.rpc.Invocation { @@ -216,5 +217,15 @@ public interface Invocation extends org.apache.dubbo.rpc.Invocation { public org.apache.dubbo.rpc.Invocation getOriginal() { return delegate; } + + @Override + public void addInvokedInvoker(org.apache.dubbo.rpc.Invoker invoker) { + delegate.addInvokedInvoker(invoker); + } + + @Override + public List> getInvokedInvokers() { + return delegate.getInvokedInvokers(); + } } } diff --git a/dubbo-compatible/src/main/java/com/alibaba/dubbo/rpc/RpcInvocation.java b/dubbo-compatible/src/main/java/com/alibaba/dubbo/rpc/RpcInvocation.java index 514dc971d5..9035294927 100644 --- a/dubbo-compatible/src/main/java/com/alibaba/dubbo/rpc/RpcInvocation.java +++ b/dubbo-compatible/src/main/java/com/alibaba/dubbo/rpc/RpcInvocation.java @@ -17,15 +17,16 @@ package com.alibaba.dubbo.rpc; -import com.alibaba.dubbo.common.Constants; -import com.alibaba.dubbo.common.URL; - import java.io.Serializable; import java.lang.reflect.Method; import java.util.Arrays; import java.util.HashMap; +import java.util.List; import java.util.Map; +import com.alibaba.dubbo.common.Constants; +import com.alibaba.dubbo.common.URL; + public class RpcInvocation implements Invocation, Serializable { private static final long serialVersionUID = -4355285085441097045L; @@ -198,6 +199,16 @@ public class RpcInvocation implements Invocation, Serializable { return value; } + @Override + public void addInvokedInvoker(org.apache.dubbo.rpc.Invoker invoker) { + throw new UnsupportedOperationException(); + } + + @Override + public List> getInvokedInvokers() { + throw new UnsupportedOperationException(); + } + @Override public String toString() { return "RpcInvocation [methodName=" + methodName + ", parameterTypes=" diff --git a/dubbo-compatible/src/test/java/org/apache/dubbo/cache/CacheTest.java b/dubbo-compatible/src/test/java/org/apache/dubbo/cache/CacheTest.java index a63ee4ddd0..b60f7bf163 100644 --- a/dubbo-compatible/src/test/java/org/apache/dubbo/cache/CacheTest.java +++ b/dubbo-compatible/src/test/java/org/apache/dubbo/cache/CacheTest.java @@ -17,17 +17,18 @@ package org.apache.dubbo.cache; +import java.util.List; +import java.util.Map; + import org.apache.dubbo.rpc.RpcInvocation; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; import com.alibaba.dubbo.cache.Cache; import com.alibaba.dubbo.cache.CacheFactory; import com.alibaba.dubbo.common.URL; import com.alibaba.dubbo.rpc.Invocation; import com.alibaba.dubbo.rpc.Invoker; -import org.junit.jupiter.api.Assertions; -import org.junit.jupiter.api.Test; - -import java.util.Map; class CacheTest { @@ -107,5 +108,15 @@ class CacheTest { public Map getAttributes() { return null; } + + @Override + public void addInvokedInvoker(org.apache.dubbo.rpc.Invoker invoker) { + + } + + @Override + public List> getInvokedInvokers() { + return null; + } } } diff --git a/dubbo-compatible/src/test/java/org/apache/dubbo/filter/LegacyInvocation.java b/dubbo-compatible/src/test/java/org/apache/dubbo/filter/LegacyInvocation.java index a2f5d8eb5c..8495a944e9 100644 --- a/dubbo-compatible/src/test/java/org/apache/dubbo/filter/LegacyInvocation.java +++ b/dubbo-compatible/src/test/java/org/apache/dubbo/filter/LegacyInvocation.java @@ -16,12 +16,13 @@ */ package org.apache.dubbo.filter; +import java.util.HashMap; +import java.util.List; +import java.util.Map; + import com.alibaba.dubbo.rpc.Invocation; import com.alibaba.dubbo.rpc.Invoker; -import java.util.HashMap; -import java.util.Map; - import static org.apache.dubbo.common.constants.CommonConstants.DUBBO_VERSION_KEY; import static org.apache.dubbo.common.constants.CommonConstants.GROUP_KEY; import static org.apache.dubbo.common.constants.CommonConstants.PATH_KEY; @@ -100,4 +101,13 @@ public class LegacyInvocation implements Invocation { return getAttachments().get(key); } + @Override + public void addInvokedInvoker(org.apache.dubbo.rpc.Invoker invoker) { + + } + + @Override + public List> getInvokedInvokers() { + return null; + } } diff --git a/dubbo-compatible/src/test/java/org/apache/dubbo/service/MockInvocation.java b/dubbo-compatible/src/test/java/org/apache/dubbo/service/MockInvocation.java index 233260f3a2..d29e339c1f 100644 --- a/dubbo-compatible/src/test/java/org/apache/dubbo/service/MockInvocation.java +++ b/dubbo-compatible/src/test/java/org/apache/dubbo/service/MockInvocation.java @@ -16,15 +16,16 @@ */ package org.apache.dubbo.service; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.function.Consumer; + import org.apache.dubbo.rpc.AttachmentsAdapter; import org.apache.dubbo.rpc.Invocation; import org.apache.dubbo.rpc.Invoker; import org.apache.dubbo.rpc.model.ServiceModel; -import java.util.HashMap; -import java.util.Map; -import java.util.function.Consumer; - import static org.apache.dubbo.common.constants.CommonConstants.DUBBO_VERSION_KEY; import static org.apache.dubbo.common.constants.CommonConstants.GROUP_KEY; import static org.apache.dubbo.common.constants.CommonConstants.PATH_KEY; @@ -179,4 +180,14 @@ public class MockInvocation implements Invocation { } return result; } + + @Override + public void addInvokedInvoker(Invoker invoker) { + + } + + @Override + public List> getInvokedInvokers() { + return null; + } } diff --git a/dubbo-rpc/dubbo-rpc-api/src/main/java/org/apache/dubbo/rpc/Invocation.java b/dubbo-rpc/dubbo-rpc-api/src/main/java/org/apache/dubbo/rpc/Invocation.java index d2d8f64365..1cb4177a90 100644 --- a/dubbo-rpc/dubbo-rpc-api/src/main/java/org/apache/dubbo/rpc/Invocation.java +++ b/dubbo-rpc/dubbo-rpc-api/src/main/java/org/apache/dubbo/rpc/Invocation.java @@ -16,15 +16,16 @@ */ package org.apache.dubbo.rpc; +import java.util.List; +import java.util.Map; +import java.util.function.Consumer; +import java.util.stream.Stream; + import org.apache.dubbo.common.Experimental; import org.apache.dubbo.rpc.model.ModuleModel; import org.apache.dubbo.rpc.model.ScopeModelUtil; import org.apache.dubbo.rpc.model.ServiceModel; -import java.util.Map; -import java.util.function.Consumer; -import java.util.stream.Stream; - /** * Invocation. (API, Prototype, NonThreadSafe) * @@ -162,4 +163,20 @@ public interface Invocation { Object get(Object key); Map getAttributes(); + + /** + * To add invoked invokers into invocation. Can be used in ClusterFilter or Filter for tracing or debugging purpose. + * Currently, only support in consumer side. + * + * @param invoker invoked invokers + */ + void addInvokedInvoker(Invoker invoker); + + /** + * Get all invoked invokers in current invocation. + * NOTICE: A curtain invoker could be invoked for twice or more if retries. + * + * @return invokers + */ + List> getInvokedInvokers(); } diff --git a/dubbo-rpc/dubbo-rpc-api/src/main/java/org/apache/dubbo/rpc/RpcInvocation.java b/dubbo-rpc/dubbo-rpc-api/src/main/java/org/apache/dubbo/rpc/RpcInvocation.java index 4ac08dbb1e..e46a6b274f 100644 --- a/dubbo-rpc/dubbo-rpc-api/src/main/java/org/apache/dubbo/rpc/RpcInvocation.java +++ b/dubbo-rpc/dubbo-rpc-api/src/main/java/org/apache/dubbo/rpc/RpcInvocation.java @@ -16,6 +16,22 @@ */ package org.apache.dubbo.rpc; +import java.io.Serializable; +import java.lang.reflect.Method; +import java.lang.reflect.Type; +import java.util.Arrays; +import java.util.Collections; +import java.util.HashMap; +import java.util.LinkedList; +import java.util.List; +import java.util.Locale; +import java.util.Map; +import java.util.concurrent.atomic.AtomicReference; +import java.util.concurrent.locks.Lock; +import java.util.concurrent.locks.ReentrantLock; +import java.util.function.Consumer; +import java.util.stream.Stream; + import org.apache.dubbo.common.URL; import org.apache.dubbo.common.utils.ReflectUtils; import org.apache.dubbo.common.utils.StringUtils; @@ -26,20 +42,6 @@ import org.apache.dubbo.rpc.model.ServiceDescriptor; import org.apache.dubbo.rpc.model.ServiceModel; import org.apache.dubbo.rpc.support.RpcUtils; -import java.io.Serializable; -import java.lang.reflect.Method; -import java.lang.reflect.Type; -import java.util.Arrays; -import java.util.Collections; -import java.util.HashMap; -import java.util.Locale; -import java.util.Map; -import java.util.concurrent.atomic.AtomicReference; -import java.util.concurrent.locks.Lock; -import java.util.concurrent.locks.ReentrantLock; -import java.util.function.Consumer; -import java.util.stream.Stream; - import static org.apache.dubbo.common.constants.CommonConstants.APPLICATION_KEY; import static org.apache.dubbo.common.constants.CommonConstants.GROUP_KEY; import static org.apache.dubbo.common.constants.CommonConstants.INTERFACE_KEY; @@ -92,6 +94,8 @@ public class RpcInvocation implements Invocation, Serializable { private transient InvokeMode invokeMode; + private transient List> invokedInvokers = new LinkedList<>(); + /** * @deprecated only for test */ @@ -311,6 +315,16 @@ public class RpcInvocation implements Invocation, Serializable { return attributes; } + @Override + public void addInvokedInvoker(Invoker invoker) { + this.invokedInvokers.add(invoker); + } + + @Override + public List> getInvokedInvokers() { + return this.invokedInvokers; + } + @Override public String getTargetServiceUniqueName() { return targetServiceUniqueName; diff --git a/dubbo-rpc/dubbo-rpc-api/src/test/java/org/apache/dubbo/rpc/RpcInvocationTest.java b/dubbo-rpc/dubbo-rpc-api/src/test/java/org/apache/dubbo/rpc/RpcInvocationTest.java index 4c226f3d0c..461867efef 100644 --- a/dubbo-rpc/dubbo-rpc-api/src/test/java/org/apache/dubbo/rpc/RpcInvocationTest.java +++ b/dubbo-rpc/dubbo-rpc-api/src/test/java/org/apache/dubbo/rpc/RpcInvocationTest.java @@ -16,10 +16,12 @@ */ package org.apache.dubbo.rpc; +import java.util.Arrays; +import java.util.HashMap; + import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; - -import java.util.HashMap; +import org.mockito.Mockito; class RpcInvocationTest { @@ -43,4 +45,20 @@ class RpcInvocationTest { invocation.setObjectAttachments(map); Assertions.assertEquals(map, invocation.getObjectAttachments()); } + + @Test + void testInvokers() { + RpcInvocation rpcInvocation = new RpcInvocation(); + + Invoker invoker1 = Mockito.mock(Invoker.class); + Invoker invoker2 = Mockito.mock(Invoker.class); + Invoker invoker3 = Mockito.mock(Invoker.class); + + rpcInvocation.addInvokedInvoker(invoker1); + rpcInvocation.addInvokedInvoker(invoker2); + rpcInvocation.addInvokedInvoker(invoker3); + rpcInvocation.addInvokedInvoker(invoker3); + + Assertions.assertEquals(Arrays.asList(invoker1, invoker2, invoker3, invoker3), rpcInvocation.getInvokedInvokers()); + } }