diff --git a/platform/platform-tests/testSrc/com/intellij/util/messages/MessageBusTest.java b/platform/platform-tests/testSrc/com/intellij/util/messages/MessageBusTest.java index 7694e8163508..546d0dcd8a54 100644 --- a/platform/platform-tests/testSrc/com/intellij/util/messages/MessageBusTest.java +++ b/platform/platform-tests/testSrc/com/intellij/util/messages/MessageBusTest.java @@ -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"); + } } diff --git a/platform/util/src/com/intellij/util/messages/impl/MessageBusImpl.java b/platform/util/src/com/intellij/util/messages/impl/MessageBusImpl.java index 44ef53835660..8291a7ba8822 100644 --- a/platform/util/src/com/intellij/util/messages/impl/MessageBusImpl.java +++ b/platform/util/src/com/intellij/util/messages/impl/MessageBusImpl.java @@ -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 map = asRoot().myWaitingBuses.get(); + final Map map = asRoot().myWaitingBuses.get(); if (map != null) { - Set buses = map.keySet(); + List buses = ContainerUtil.filter(map.keySet(), new Condition() { + @Override + public boolean value(MessageBusImpl bus) { + return ensureAlive(map, bus); + } + }); if (!buses.isEmpty()) { - pumpWaitingBuses(map, new ArrayList(buses)); + pumpWaitingBuses(buses); } } } } - private static void pumpWaitingBuses(Map map, ArrayList buses) { + private static void pumpWaitingBuses(List buses) { List exceptions = null; for (MessageBusImpl bus : buses) { - if (!ensureAlive(map, bus)) continue; + if (bus.myDisposed) continue; exceptions = appendExceptions(exceptions, bus.doPumpMessages()); }