message bus - use method handles to invoke method

GitOrigin-RevId: 4eac0d10cdf62cdd809bc4f31262e3ced9b0b82f
This commit is contained in:
Vladimir Krivosheev
2020-11-20 07:55:45 +00:00
committed by intellij-monorepo-bot
parent e7aba6f006
commit 4f6ce9c63c
8 changed files with 92 additions and 34 deletions
@@ -6,13 +6,14 @@ import com.intellij.util.messages.Topic;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
import java.lang.reflect.Method;
import java.lang.invoke.MethodHandle;
import java.util.Arrays;
import java.util.List;
final class Message<L> {
final Topic<L> topic;
final Method listenerMethod;
final String methodName;
final MethodHandle method;
final Object[] args;
final List<L> handlers;
final @Nullable ClientId clientId;
@@ -21,10 +22,10 @@ final class Message<L> {
// see note about pumpMessages in createPublisher (invoking job handlers can be stopped and continued as part of another pumpMessages call)
int currentHandlerIndex;
Message(@NotNull Topic<L> topic, @NotNull Method listenerMethod, Object[] args, @NotNull List<L> handlers) {
Message(@NotNull Topic<L> topic, @NotNull MethodHandle method, @NotNull String methodName, Object[] args, @NotNull List<L> handlers) {
this.topic = topic;
listenerMethod.setAccessible(true);
this.listenerMethod = listenerMethod;
this.method = method;
this.methodName = methodName;
this.args = args;
this.handlers = handlers;
clientId = ClientId.getCurrentOrNull();
@@ -34,7 +35,7 @@ final class Message<L> {
public String toString() {
return "Message(" +
"topic=" + topic +
", listenerMethod=" + listenerMethod +
", method=" + methodName +
", args=" + Arrays.toString(args) +
", handlers=" + handlers +
')';
@@ -16,6 +16,7 @@ import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
import org.jetbrains.annotations.TestOnly;
import java.lang.invoke.MethodHandle;
import java.lang.reflect.InvocationHandler;
import java.lang.reflect.Method;
import java.lang.reflect.Proxy;
@@ -127,13 +128,10 @@ public class MessageBusImpl implements MessageBus {
public final @NotNull <L> L syncPublisher(@NotNull Topic<L> topic) {
checkNotDisposed();
//noinspection unchecked
return (L)publisherCache.computeIfAbsent(topic, this::createPublisherInvocationHandler);
}
private @NotNull <L> L createPublisherInvocationHandler(@NotNull Topic<L> topic) {
Class<L> listenerClass = topic.getListenerClass();
//noinspection unchecked
return (L)Proxy.newProxyInstance(listenerClass.getClassLoader(), new Class[]{listenerClass}, createPublisher(topic, topic.getBroadcastDirection()));
return (L)publisherCache.computeIfAbsent(topic, topic1 -> {
Class<?> listenerClass = topic1.getListenerClass();
return Proxy.newProxyInstance(listenerClass.getClassLoader(), new Class[]{listenerClass}, createPublisher(topic1, topic1.getBroadcastDirection()));
});
}
@NotNull
@@ -163,7 +161,7 @@ public class MessageBusImpl implements MessageBus {
@Override
public final Object invoke(Object proxy, Method method, Object[] args) {
if (method.getDeclaringClass().getName().equals("java.lang.Object")) {
if (method.getDeclaringClass() == Object.class) {
return EventDispatcher.handleObjectMethod(proxy, args, method.getName());
}
@@ -213,14 +211,14 @@ public class MessageBusImpl implements MessageBus {
@Nullable JobQueue jobQueue,
@Nullable MessageDeliveryListener messageDeliveryListener,
@Nullable List<Throwable> exceptions) {
MethodHandle methodHandle = MethodHandleCache.compute(method, args);
if (jobQueue == null) {
for (L handler : handlers) {
exceptions = invokeListener(method, args, topic, handler, messageDeliveryListener, exceptions);
exceptions = invokeListener(methodHandle, method.getName(), args, topic, handler, messageDeliveryListener, exceptions);
}
}
else {
Message<L> message = new Message<>(topic, method, args, handlers);
jobQueue.queue.offerLast(message);
jobQueue.queue.offerLast(new Message<>(topic, methodHandle, method.getName(), args, handlers));
}
return exceptions;
}
@@ -423,7 +421,7 @@ public class MessageBusImpl implements MessageBus {
}
job.currentHandlerIndex++;
exceptions = invokeListener(job.listenerMethod, job.args, job.topic, handlers.get(index), messageDeliveryListener, exceptions);
exceptions = invokeListener(job.method, job.methodName, job.args, job.topic, handlers.get(index), messageDeliveryListener, exceptions);
if (++index != job.currentHandlerIndex) {
// handler published some event and message queue including current job was processed as result, so, stop processing
return exceptions;
@@ -620,7 +618,7 @@ public class MessageBusImpl implements MessageBus {
job.handlers.addAll(connectionHandlers);
}
else {
filteredJob = new Message<>(job.topic, job.listenerMethod, job.args, connectionHandlers);
filteredJob = new Message<>(job.topic, job.method, job.methodName, job.args, connectionHandlers);
}
if (newJobs == null) {
newJobs = new SmartList<>();
@@ -638,7 +636,8 @@ public class MessageBusImpl implements MessageBus {
}
// args is not null
private static @Nullable <L> List<Throwable> invokeListener(@NotNull Method method,
private static @Nullable <L> List<Throwable> invokeListener(@NotNull MethodHandle methodHandle,
@NotNull String methodName,
Object[] args,
@NotNull Topic<L> topic,
@NotNull L handler,
@@ -646,23 +645,38 @@ public class MessageBusImpl implements MessageBus {
@Nullable List<Throwable> exceptions) {
try {
if (handler instanceof MessageHandler) {
((MessageHandler)handler).handle(method, args);
((MessageHandler)handler).handle(methodHandle, args);
}
else if (messageDeliveryListener == null) {
method.invoke(handler, args);
invokeMethod(handler, args, methodHandle);
}
else {
long startTime = System.nanoTime();
method.invoke(handler, args);
messageDeliveryListener.messageDelivered(topic, method.getName(), handler, System.nanoTime() - startTime);
invokeMethod(handler, args, methodHandle);
messageDeliveryListener.messageDelivered(topic, methodName, handler, System.nanoTime() - startTime);
}
}
catch (AbstractMethodError e) {
// do nothing for AbstractMethodError. This listener just does not implement something newly added yet.
}
catch (Throwable e) {
exceptions = EventDispatcher.handleException(e, exceptions);
if (exceptions == null) {
exceptions = new ArrayList<>();
}
exceptions.add(e);
}
return exceptions;
}
private static void invokeMethod(@NotNull Object handler, Object[] args, MethodHandle methodHandle) throws Throwable {
if (args == null) {
methodHandle.invoke(handler);
}
else {
methodHandle.bindTo(handler).invokeExact(args);
}
}
protected void disconnectPluginConnections(@NotNull Predicate<? super Class<?>> predicate) {
for (MessageHandlerHolder holder : subscribers) {
holder.disconnectIfNeeded(predicate);
@@ -0,0 +1,39 @@
// Copyright 2000-2020 JetBrains s.r.o. Use of this source code is governed by the Apache 2.0 license that can be found in the LICENSE file.
package com.intellij.util.messages.impl;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
import java.lang.invoke.MethodHandle;
import java.lang.invoke.MethodHandles;
import java.lang.reflect.Method;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
final class MethodHandleCache {
private static final MethodHandles.Lookup LOOKUP = MethodHandles.lookup();
private static final ClassValue<ConcurrentMap<Method, MethodHandle>> CACHE = new ConcurrentMapClassValue();
static @NotNull MethodHandle compute(@NotNull Method method, Object @Nullable [] args) {
// method name cannot be used as key because for one class maybe several methods with the same name and different set of parameters
return CACHE.get(method.getDeclaringClass()).computeIfAbsent(method, method1 -> {
method1.setAccessible(true);
MethodHandle result;
try {
result = LOOKUP.unreflect(method1);
}
catch (IllegalAccessException e) {
throw new RuntimeException(e);
}
return args == null ? result : result.asSpreader(Object[].class, args.length);
});
}
// as static to ensure that class doesn't reference anything else (otherwise maybe memory leak)
private static final class ConcurrentMapClassValue extends ClassValue<ConcurrentMap<Method, MethodHandle>> {
@Override
protected ConcurrentMap<Method, MethodHandle> computeValue(@NotNull Class<?> type) {
return new ConcurrentHashMap<>(8);
}
}
}
@@ -12,11 +12,11 @@ import org.jetbrains.annotations.Nullable;
*/
public interface MessageBusConnection extends SimpleMessageBusConnection, Disposable {
/**
* Subscribes to the target topic within the current connection using {@link #setDefaultHandler(MessageHandler) default handler}.
* Subscribes to the target topic within the current connection using {@link #setDefaultHandler(Runnable) default handler}.
*
* @param topic target endpoint
* @param <L> interface for working with the target topic
* @throws IllegalStateException if {@link #setDefaultHandler(MessageHandler) default handler} hasn't been defined or
* @throws IllegalStateException if {@link #setDefaultHandler(Runnable) default handler} hasn't been defined or
* has incompatible type with the {@link Topic#getListenerClass() topic's business interface}
* or if target topic is already subscribed within the current connection
*/
@@ -29,6 +29,10 @@ public interface MessageBusConnection extends SimpleMessageBusConnection, Dispos
*/
void setDefaultHandler(@Nullable MessageHandler handler);
default void setDefaultHandler(@NotNull Runnable runnable) {
setDefaultHandler((event, params) -> runnable.run());
}
/**
* Forces to process any queued but not delivered events.
*
@@ -3,7 +3,7 @@ package com.intellij.util.messages;
import org.jetbrains.annotations.NotNull;
import java.lang.reflect.Method;
import java.lang.invoke.MethodHandle;
/**
* Defines contract for generic messages subscriber processor.
@@ -16,5 +16,5 @@ public interface MessageHandler {
* @param event information about target method called by the publisher
* @param params called method arguments
*/
void handle(@NotNull Method event, Object... params);
void handle(@NotNull MethodHandle event, Object... params);
}
@@ -134,7 +134,7 @@ public final class DocRenderItem {
else {
if (editor.getUserData(LISTENERS_DISPOSABLE) == null) {
MessageBusConnection connection = ApplicationManager.getApplication().getMessageBus().connect();
connection.setDefaultHandler((event, params) -> updateInlays(editor, true));
connection.setDefaultHandler(() -> updateInlays(editor, true));
connection.subscribe(EditorColorsManager.TOPIC);
connection.subscribe(LafManagerListener.TOPIC);
EditorFactory.getInstance().addEditorFactoryListener(new EditorFactoryListener() {
@@ -46,7 +46,7 @@ public class ModuleManagerComponent extends ModuleManagerImpl {
super(project);
myMessageBusConnection = project.getMessageBus().connect(this);
myMessageBusConnection.setDefaultHandler((event, params) -> cleanCachedStuff());
myMessageBusConnection.setDefaultHandler(() -> cleanCachedStuff());
myMessageBusConnection.subscribe(ProjectTopics.PROJECT_ROOTS);
// default project doesn't have modules
@@ -284,13 +284,13 @@ public class MessageBusTest implements MessageBusOwner {
public void t12() {
}
});
for (int i = 0; i < 1000; i++) {
for (int i = 0; i < 1_000; i++) {
new MessageBusImpl(this, childBus);
}
PlatformTestUtil.assertTiming("Too long", 3000, () -> {
PlatformTestUtil.assertTiming("Too long", 3_000, () -> {
T1Listener publisher = myBus.syncPublisher(TOPIC1);
for (int i = 0; i < 1000000; i++) {
for (int i = 0; i < 1_000_000; i++) {
publisher.t11();
}
});