lazy message bus listeners on project level — part 1

GitOrigin-RevId: 2e27e3a04e356a4dc3647073760025d8c7f4c552
This commit is contained in:
Vladimir Krivosheev
2019-08-01 14:03:25 +03:00
committed by intellij-monorepo-bot
parent 9fdb29cff1
commit a54027f8a4
6 changed files with 46 additions and 77 deletions
@@ -15,6 +15,6 @@ public class MessageBusFactoryImpl extends MessageBusFactory {
@NotNull
@Override
public MessageBus createMessageBus(@NotNull Object owner, @NotNull MessageBus parentBus) {
return new MessageBusImpl(owner, parentBus);
return new MessageBusImpl(owner, (MessageBusImpl)parentBus);
}
}
@@ -65,17 +65,21 @@ public class MessageBusImpl implements MessageBus {
private final Disposable myConnectionDisposable;
private MessageDeliveryListener myMessageDeliveryListener;
public MessageBusImpl(@NotNull Object owner, @NotNull MessageBus parentBus) {
private final MessageBusConnectionImpl myLazyConnection;
public MessageBusImpl(@NotNull Object owner, @NotNull MessageBusImpl parentBus) {
myOwner = owner + " of " + owner.getClass();
myConnectionDisposable = Disposer.newDisposable(myOwner);
myParentBus = (MessageBusImpl)parentBus;
myRootBus = myParentBus.myRootBus;
synchronized (myParentBus.myChildBuses) {
myOrder = myParentBus.nextOrder();
myParentBus.myChildBuses.add(this);
myParentBus = parentBus;
myRootBus = parentBus.myRootBus;
synchronized (parentBus.myChildBuses) {
myOrder = parentBus.nextOrder();
parentBus.myChildBuses.add(this);
}
LOG.assertTrue(myParentBus.myChildBuses.contains(this));
LOG.assertTrue(parentBus.myChildBuses.contains(this));
myRootBus.clearSubscriberCache();
// only for project
myLazyConnection = parentBus.myParentBus == null ? connect() : null;
}
// root message bus constructor
@@ -84,6 +88,7 @@ public class MessageBusImpl implements MessageBus {
myConnectionDisposable = Disposer.newDisposable(myOwner);
myOrder = ArrayUtil.EMPTY_INT_ARRAY;
myRootBus = (RootBus)this;
myLazyConnection = connect();
}
/**
@@ -130,7 +135,7 @@ public class MessageBusImpl implements MessageBus {
LOG.assertTrue(removed);
}
private static class DeliveryJob {
private static final class DeliveryJob {
DeliveryJob(@NotNull MessageBusConnectionImpl connection, @NotNull Message message) {
this.connection = connection;
this.message = message;
@@ -161,11 +166,6 @@ public class MessageBusImpl implements MessageBus {
return connection;
}
@NotNull
protected MessageBusConnectionImpl createConnectionForLazyListeners() {
return connect();
}
@Override
@NotNull
public <L> L syncPublisher(@NotNull Topic<L> topic) {
@@ -204,7 +204,6 @@ public class MessageBusImpl implements MessageBus {
List<ListenerDescriptor> listenerDescriptors = myTopicClassToListenerClass.remove(listenerClass.getName());
if (listenerDescriptors != null) {
MessageBusConnectionImpl connection = createConnectionForLazyListeners();
List<Object> listeners = new SmartList<>();
for (ListenerDescriptor listenerDescriptor : listenerDescriptors) {
ClassLoader classLoader = listenerDescriptor.pluginDescriptor.getPluginClassLoader();
@@ -217,7 +216,7 @@ public class MessageBusImpl implements MessageBus {
}
if (!listeners.isEmpty()) {
connection.subscribe(topic, listeners);
myLazyConnection.subscribe(topic, listeners);
}
}
@@ -531,16 +530,8 @@ public class MessageBusImpl implements MessageBus {
*/
private final ThreadLocal<SortedMap<MessageBusImpl, Integer>> myWaitingBuses = new ThreadLocal<>();
private final MessageBusConnectionImpl myLazyConnection = connect();
private volatile boolean myClearedSubscribersCache;
@NotNull
@Override
protected MessageBusConnectionImpl createConnectionForLazyListeners() {
return myLazyConnection;
}
@Override
void clearSubscriberCache() {
if (myClearedSubscribersCache) return;
@@ -1,18 +1,4 @@
/*
* Copyright 2000-2016 JetBrains s.r.o.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
// Copyright 2000-2019 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.mock;
import com.intellij.openapi.Disposable;
@@ -42,7 +28,7 @@ public class MockComponentManager extends UserDataHolderBase implements Componen
private final MessageBus myMessageBus = new MessageBusFactoryImpl().createMessageBus(this);
private final MutablePicoContainer myPicoContainer;
private final Map<Class, Object> myComponents = new HashMap<>();
private final Map<Class<?>, Object> myComponents = new HashMap<>();
private final Set<Object> myDisposableComponents = ContainerUtil.newConcurrentSet();
private boolean myDisposed;
@@ -40,8 +40,12 @@ import com.intellij.ui.CustomProtocolHandler
import com.intellij.ui.mac.MacOSApplicationProvider
import com.intellij.ui.mac.touchbar.TouchBarsManager
import com.intellij.util.ArrayUtilRt
import com.intellij.util.SmartList
import com.intellij.util.concurrency.AppExecutorUtil
import com.intellij.util.containers.ContainerUtil
import com.intellij.util.io.exists
import com.intellij.util.messages.ListenerDescriptor
import com.intellij.util.messages.impl.MessageBusImpl
import com.intellij.util.ui.AsyncProcessIcon
import com.intellij.util.ui.accessibility.ScreenReader
import net.miginfocom.layout.PlatformDefaults
@@ -184,7 +188,7 @@ fun registerRegistryAndContainerAndInitStore(pluginDescriptorsFuture: Completabl
app.registerComponents(pluginDescriptors)
}
initAppActivity.runChild("add message bus listeners") {
app.registerMessageBusListeners(pluginDescriptors, false)
registerMessageBusListeners(pluginDescriptors, app)
}
// yes, at this moment initSystemProperties or RegistryKeyBean.addKeysFromPlugins maybe not yet performed, but it doesn't affect because not used.
@@ -525,4 +529,25 @@ fun preloadServices(app: ApplicationImpl): CompletableFuture<Void?> {
}
}, appExecutorService)
})
}
private fun registerMessageBusListeners(pluginDescriptors: List<IdeaPluginDescriptor>, app: ApplicationImpl) {
val map = ContainerUtil.newConcurrentMap<String, MutableList<ListenerDescriptor>>()
val isHeadlessMode = app.isHeadlessEnvironment
val isUnitTestMode = app.isUnitTestMode
for (descriptor in pluginDescriptors) {
val listeners = (descriptor as IdeaPluginDescriptorImpl).app.listeners
for (listener in listeners) {
if (isUnitTestMode && !listener.activeInTestMode) {
continue
}
if (isHeadlessMode && !listener.activeInHeadlessMode) {
continue
}
map.getOrPut(listener.topicClassName) { SmartList() }.add(listener)
}
}
(app.messageBus as MessageBusImpl).setLazyListeners(map)
}
@@ -47,12 +47,9 @@ import com.intellij.util.*;
import com.intellij.util.concurrency.AppExecutorUtil;
import com.intellij.util.concurrency.AppScheduledExecutorService;
import com.intellij.util.concurrency.Semaphore;
import com.intellij.util.containers.ContainerUtil;
import com.intellij.util.containers.Stack;
import com.intellij.util.io.storage.HeavyProcessLatch;
import com.intellij.util.messages.ListenerDescriptor;
import com.intellij.util.messages.Topic;
import com.intellij.util.messages.impl.MessageBusImpl;
import com.intellij.util.ui.EdtInvocationManager;
import com.intellij.util.ui.UIUtil;
import org.jetbrains.annotations.*;
@@ -170,35 +167,6 @@ public class ApplicationImpl extends PlatformComponentManagerImpl implements App
IdeEventQueue.getInstance();
}
// this method is not in ApplicationImpl constructor because application starter can perform this activity in parallel to another task
@ApiStatus.Internal
public void registerMessageBusListeners(@NotNull List<? extends IdeaPluginDescriptor> pluginDescriptors, boolean isUnitTestMode) {
ConcurrentMap<String, List<ListenerDescriptor>> map = ContainerUtil.newConcurrentMap();
boolean isHeadlessMode = isHeadlessEnvironment();
for (IdeaPluginDescriptor descriptor : pluginDescriptors) {
List<ListenerDescriptor> listeners = ((IdeaPluginDescriptorImpl)descriptor).getApp().getListeners();
if (!listeners.isEmpty()) {
for (ListenerDescriptor listener : listeners) {
if (isUnitTestMode && !listener.activeInTestMode) {
continue;
}
if (isHeadlessMode && !listener.activeInHeadlessMode) {
continue;
}
List<ListenerDescriptor> list = map.get(listener.topicClassName);
if (list == null) {
list = new SmartList<>();
map.put(listener.topicClassName, list);
}
list.add(listener);
}
}
}
((MessageBusImpl)getMessageBus()).setLazyListeners(map);
}
/**
* Executes a {@code runnable} in an "impatient" mode.
* In this mode any attempt to call {@link #runReadAction(Runnable)}
@@ -25,7 +25,7 @@ import java.util.concurrent.CountDownLatch;
import java.util.concurrent.atomic.AtomicReference;
public class MessageBusTest extends LightPlatformTestCase {
private MessageBus myBus;
private MessageBusImpl myBus;
private List<String> myLog;
public interface T1Listener {
@@ -77,12 +77,11 @@ public class MessageBusTest extends LightPlatformTestCase {
}
}
private Disposable myParentDisposable = Disposer.newDisposable();
@Override
protected void setUp() throws Exception {
super.setUp();
myBus = MessageBusFactory.newMessageBus(this);
myBus = (MessageBusImpl)MessageBusFactory.newMessageBus(this);
Disposer.register(myParentDisposable, myBus);
myLog = new ArrayList<>();
}
@@ -269,7 +268,7 @@ public class MessageBusTest extends LightPlatformTestCase {
final int threadsNumber = 10;
final AtomicReference<Throwable> exception = new AtomicReference<>();
final CountDownLatch latch = new CountDownLatch(threadsNumber);
final MessageBus parentBus = MessageBusFactory.newMessageBus("parent");
MessageBusImpl parentBus = (MessageBusImpl)MessageBusFactory.newMessageBus("parent");
Disposer.register(myParentDisposable, parentBus);
List<Thread> threads = new ArrayList<>();
final int iterationsNumber = 100;