allow to dispose message bus during event processing when it's in the queue

to support cases like CidrProjectFixture.closeProjectAndCleanup being called from projectOpened listener
This commit is contained in:
peter
2017-05-15 16:16:19 +02:00
parent 3dc7fe8e98
commit d8f61ec1ac
2 changed files with 41 additions and 5 deletions
@@ -335,4 +335,34 @@ public class MessageBusTest extends TestCase {
childBus.connect().subscribe(RUNNABLE_TOPIC, () -> assertFalse(myBus.hasUndeliveredEvents(RUNNABLE_TOPIC)));
myBus.syncPublisher(RUNNABLE_TOPIC).run();
}
public void testDisposingBusInsideEvent() {
MessageBusImpl child = new MessageBusImpl(this, myBus);
myBus.connect().subscribe(TOPIC1, new T1Listener() {
@Override
public void t11() {
myLog.add("root 11");
myBus.syncPublisher(TOPIC1).t12();
child.dispose();
}
@Override
public void t12() {
myLog.add("root 12");
}
});
child.connect().subscribe(TOPIC1, new T1Listener() {
@Override
public void t11() {
myLog.add("child 11");
}
@Override
public void t12() {
myLog.add("child 12");
}
});
myBus.syncPublisher(TOPIC1).t11();
assertEvents("root 11", "child 11", "root 12", "child 12");
}
}
@@ -18,6 +18,7 @@ package com.intellij.util.messages.impl;
import com.intellij.openapi.Disposable;
import com.intellij.openapi.diagnostic.Logger;
import com.intellij.openapi.progress.ProcessCanceledException;
import com.intellij.openapi.util.Condition;
import com.intellij.openapi.util.Disposer;
import com.intellij.util.ConcurrencyUtil;
import com.intellij.util.SmartList;
@@ -380,20 +381,25 @@ public class MessageBusImpl implements MessageBus {
myParentBus.pumpMessages();
}
else {
Map<MessageBusImpl, Integer> map = asRoot().myWaitingBuses.get();
final Map<MessageBusImpl, Integer> map = asRoot().myWaitingBuses.get();
if (map != null) {
Set<MessageBusImpl> buses = map.keySet();
List<MessageBusImpl> buses = ContainerUtil.filter(map.keySet(), new Condition<MessageBusImpl>() {
@Override
public boolean value(MessageBusImpl bus) {
return ensureAlive(map, bus);
}
});
if (!buses.isEmpty()) {
pumpWaitingBuses(map, new ArrayList<MessageBusImpl>(buses));
pumpWaitingBuses(buses);
}
}
}
}
private static void pumpWaitingBuses(Map<MessageBusImpl, Integer> map, ArrayList<MessageBusImpl> buses) {
private static void pumpWaitingBuses(List<MessageBusImpl> buses) {
List<Throwable> exceptions = null;
for (MessageBusImpl bus : buses) {
if (!ensureAlive(map, bus)) continue;
if (bus.myDisposed) continue;
exceptions = appendExceptions(exceptions, bus.doPumpMessages());
}