From 13c01efd12340f1bf065c6b1ca329c63e384bf65 Mon Sep 17 00:00:00 2001 From: "ken.lj" Date: Wed, 26 Jun 2019 15:56:20 +0800 Subject: [PATCH] test register callback once --- .../rpc/protocol/ProtocolFilterWrapper.java | 72 +++++++++++++++---- 1 file changed, 59 insertions(+), 13 deletions(-) diff --git a/dubbo-rpc/dubbo-rpc-api/src/main/java/org/apache/dubbo/rpc/protocol/ProtocolFilterWrapper.java b/dubbo-rpc/dubbo-rpc-api/src/main/java/org/apache/dubbo/rpc/protocol/ProtocolFilterWrapper.java index bfaf14d2a5..f54d07688b 100644 --- a/dubbo-rpc/dubbo-rpc-api/src/main/java/org/apache/dubbo/rpc/protocol/ProtocolFilterWrapper.java +++ b/dubbo-rpc/dubbo-rpc-api/src/main/java/org/apache/dubbo/rpc/protocol/ProtocolFilterWrapper.java @@ -48,9 +48,12 @@ public class ProtocolFilterWrapper implements Protocol { this.protocol = protocol; } + + private static Invoker buildInvokerChain(final Invoker invoker, String key, String group) { Invoker last = invoker; List filters = ExtensionLoader.getExtensionLoader(Filter.class).getActivateExtension(invoker.getUrl(), key, group); + if (!filters.isEmpty()) { for (int i = filters.size() - 1; i >= 0; i--) { final Filter filter = filters.get(i); @@ -87,18 +90,7 @@ public class ProtocolFilterWrapper implements Protocol { } throw e; } - return asyncResult.thenApplyWithContext(r -> { - // onResponse callback - if (filter instanceof ListenableFilter) { - Filter.Listener listener = ((ListenableFilter) filter).listener(); - if (listener != null) { - listener.onResponse(r, invoker, invocation); - } - } else { - filter.onResponse(r, invoker, invocation); - } - return r; - }); + return asyncResult; } @Override @@ -113,7 +105,8 @@ public class ProtocolFilterWrapper implements Protocol { }; } } - return last; + + return new CallbackRegistrationInvoker<>(last, filters); } @Override @@ -142,4 +135,57 @@ public class ProtocolFilterWrapper implements Protocol { protocol.destroy(); } + static class CallbackRegistrationInvoker implements Invoker { + + private final Invoker filterInvoker; + private final List filters; + + public CallbackRegistrationInvoker(Invoker filterInvoker, List filters) { + this.filterInvoker = filterInvoker; + this.filters = filters; + } + + @Override + public Result invoke(Invocation invocation) throws RpcException { + Result asyncResult = filterInvoker.invoke(invocation); + + asyncResult.thenApplyWithContext(r -> { + for (int i = filters.size() - 1; i >= 0; i--) { + Filter filter = filters.get(i); + // onResponse callback + if (filter instanceof ListenableFilter) { + Filter.Listener listener = ((ListenableFilter) filter).listener(); + if (listener != null) { + listener.onResponse(r, filterInvoker, invocation); + } + } else { + filter.onResponse(r, filterInvoker, invocation); + } + } + return r; + }); + + return asyncResult; + } + + @Override + public Class getInterface() { + return filterInvoker.getInterface(); + } + + @Override + public URL getUrl() { + return filterInvoker.getUrl(); + } + + @Override + public boolean isAvailable() { + return filterInvoker.isAvailable(); + } + + @Override + public void destroy() { + filterInvoker.destroy(); + } + } }