warmup: flush all of MergingUpdateQueue-s

GitOrigin-RevId: b63ca37d98a9f917d1690b6cfb9e184b7b5bfea1
This commit is contained in:
Dmitry Batkovich
2022-11-15 08:14:19 +00:00
committed by intellij-monorepo-bot
parent e431dbb226
commit b5ba09d72b
2 changed files with 59 additions and 16 deletions
@@ -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<Map<Update, Update>> myScheduledUpdates = ConcurrentCollectionFactory.createConcurrentIntObjectMap();
private static final Set<MergingUpdateQueue> 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() {
@@ -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))
}
}