diff --git a/platform/ide-core/src/com/intellij/util/ui/update/MergingUpdateQueue.java b/platform/ide-core/src/com/intellij/util/ui/update/MergingUpdateQueue.java index 6b0c85e15f09..fe9606606167 100644 --- a/platform/ide-core/src/com/intellij/util/ui/update/MergingUpdateQueue.java +++ b/platform/ide-core/src/com/intellij/util/ui/update/MergingUpdateQueue.java @@ -12,6 +12,7 @@ import com.intellij.openapi.progress.ProcessCanceledException; import com.intellij.openapi.util.Disposer; import com.intellij.util.Alarm; import com.intellij.util.AlarmFactory; +import com.intellij.util.SystemProperties; import com.intellij.util.containers.ConcurrentIntObjectMap; import com.intellij.util.containers.ContainerUtil; import com.intellij.util.ui.EdtInvocationManager; @@ -37,6 +38,10 @@ public class MergingUpdateQueue implements Runnable, Disposable, Activatable { private volatile boolean mySuspended; private final ConcurrentIntObjectMap> myScheduledUpdates = ConcurrentCollectionFactory.createConcurrentIntObjectMap(); + private static final Set ourQueues = + SystemProperties.getBooleanProperty("intellij.MergingUpdateQueue.enable.global.flusher", false) + ? ContainerUtil.newConcurrentSet() + : null; private final Alarm myWaiterForMerge; @@ -122,6 +127,10 @@ public class MergingUpdateQueue implements Runnable, Disposable, Activatable { UiNotifyConnector connector = new UiNotifyConnector(activationComponent, this); Disposer.register(this, connector); } + + if (ourQueues != null) { + ourQueues.add(this); + } } public void setMergingTimeSpan(int timeSpan) { @@ -238,6 +247,15 @@ public class MergingUpdateQueue implements Runnable, Disposable, Activatable { flush(); } + @ApiStatus.Internal + public static void flushAllQueues() { + if (ourQueues != null) { + for (MergingUpdateQueue queue : ourQueues) { + queue.flush(); + } + } + } + /** * executes all scheduled requests in the current thread. * Please note that requests that started execution before this method call are not waited for completion. @@ -245,7 +263,7 @@ public class MergingUpdateQueue implements Runnable, Disposable, Activatable { public void flush() { synchronized (myScheduledUpdates) { if (myScheduledUpdates.isEmpty()) { - finishActivity(); + //finishActivity(); return; } } @@ -413,10 +431,17 @@ public class MergingUpdateQueue implements Runnable, Disposable, Activatable { @Override public void dispose() { - myDisposed = true; - myActive = false; - finishActivity(); - clearWaiter(); + try { + myDisposed = true; + myActive = false; + finishActivity(); + clearWaiter(); + } + finally { + if (ourQueues != null) { + ourQueues.remove(this); + } + } } private void clearWaiter() { diff --git a/platform/warmup/src/com/intellij/warmup/warmupUtils.kt b/platform/warmup/src/com/intellij/warmup/warmupUtils.kt index 7c4565fde63b..9cf5f41826fe 100644 --- a/platform/warmup/src/com/intellij/warmup/warmupUtils.kt +++ b/platform/warmup/src/com/intellij/warmup/warmupUtils.kt @@ -1,12 +1,17 @@ // Copyright 2000-2022 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license. package com.intellij.warmup +import com.intellij.openapi.application.EDT import com.intellij.openapi.progress.impl.CoreProgressManager import com.intellij.util.indexing.FileBasedIndex import com.intellij.util.indexing.FileBasedIndexEx +import com.intellij.util.ui.update.MergingUpdateQueue import com.intellij.warmup.util.ConsoleLog import com.intellij.warmup.util.runTaskAndLogTime +import com.intellij.warmup.util.yieldThroughInvokeLater +import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.time.delay +import kotlinx.coroutines.withContext import java.time.Duration import kotlin.system.exitProcess @@ -26,18 +31,31 @@ suspend fun waitUntilProgressTasksAreFinishedOrFail() { private suspend fun waitUntilProgressTasksAreFinished() { runTaskAndLogTime("Awaiting for progress tasks") { - val timeout = System.getProperty("ide.progress.tasks.awaiting.timeout.min", "60").toLongOrNull() ?: 60 - val startTime = System.currentTimeMillis() - while (CoreProgressManager.getCurrentIndicators().isNotEmpty()) { - if (System.currentTimeMillis() - startTime > Duration.ofMinutes(timeout).toMillis()) { - val timeoutMessage = StringBuilder("Progress tasks awaiting timeout.\n") - timeoutMessage.appendLine("Not finished tasks:") - for (indicator in CoreProgressManager.getCurrentIndicators()) { - timeoutMessage.appendLine(" - ${indicator.text}") - } - error(timeoutMessage) + while (true) { + withContext(Dispatchers.EDT) { + MergingUpdateQueue.flushAllQueues() } - delay(Duration.ofMillis(1000)) + yieldThroughInvokeLater() + if (CoreProgressManager.getCurrentIndicators().isEmpty()) { + return@runTaskAndLogTime + } + waitCurrentProgressIndicators() } } +} + +private suspend fun waitCurrentProgressIndicators() { + val timeout = System.getProperty("ide.progress.tasks.awaiting.timeout.min", "60").toLongOrNull() ?: 60 + val startTime = System.currentTimeMillis() + while (CoreProgressManager.getCurrentIndicators().isNotEmpty()) { + if (System.currentTimeMillis() - startTime > Duration.ofMinutes(timeout).toMillis()) { + val timeoutMessage = StringBuilder("Progress tasks awaiting timeout.\n") + timeoutMessage.appendLine("Not finished tasks:") + for (indicator in CoreProgressManager.getCurrentIndicators()) { + timeoutMessage.appendLine(" - ${indicator.text}") + } + error(timeoutMessage) + } + delay(Duration.ofMillis(1000)) + } } \ No newline at end of file