From 1a16205a961602fb9316cb0b5b325e78f9e44a6b Mon Sep 17 00:00:00 2001 From: Anton Makeev Date: Tue, 13 Jun 2017 12:52:58 +0200 Subject: [PATCH] QueueProcessor.waitFor should return correctly when queue is shut down --- .../util/concurrency/QueueProcessor.java | 54 +++++++++++++------ .../concurrency}/BackgroundTaskQueueTest.java | 2 +- .../util/concurrency/QueueProcessorTest.kt | 33 ++++++++++++ 3 files changed, 72 insertions(+), 17 deletions(-) rename platform/{vcs-tests/testSrc/com/intellij/openapi/vcs/changes/committed => platform-tests/testSrc/com/intellij/util/concurrency}/BackgroundTaskQueueTest.java (99%) create mode 100644 platform/platform-tests/testSrc/com/intellij/util/concurrency/QueueProcessorTest.kt diff --git a/platform/platform-api/src/com/intellij/util/concurrency/QueueProcessor.java b/platform/platform-api/src/com/intellij/util/concurrency/QueueProcessor.java index a0543a51f343..d86860f0226d 100644 --- a/platform/platform-api/src/com/intellij/util/concurrency/QueueProcessor.java +++ b/platform/platform-api/src/com/intellij/util/concurrency/QueueProcessor.java @@ -51,20 +51,6 @@ public class QueueProcessor { private final PairConsumer myProcessor; private final Deque myQueue = new ArrayDeque<>(); - private final Runnable myContinuationContext = new Runnable() { - @Override - public void run() { - synchronized (myQueue) { - isProcessing = false; - if (myQueue.isEmpty()) { - myQueue.notifyAll(); - } - else { - startProcessing(); - } - } - } - }; private boolean isProcessing; private boolean myStarted; @@ -144,6 +130,18 @@ public class QueueProcessor { } } } + + private void finishProcessing(boolean continueProcessing) { + synchronized (myQueue) { + isProcessing = false; + if (myQueue.isEmpty()) { + myQueue.notifyAll(); + } + else if (continueProcessing){ + startProcessing(); + } + } + } public void add(@NotNull T t, ModalityState state) { synchronized (myQueue) { @@ -190,6 +188,27 @@ public class QueueProcessor { } } } + + public boolean waitFor(long timeoutMS) { + synchronized (myQueue) { + long start = System.currentTimeMillis(); + + while (isProcessing) { + long rest = timeoutMS - (System.currentTimeMillis() - start); + + if (rest <= 0) return !isProcessing; + + try { + myQueue.wait(rest); + } + catch (InterruptedException e) { + //ok + } + } + + return true; + } + } private boolean startProcessing() { LOG.assertTrue(Thread.holdsLock(myQueue)); @@ -200,8 +219,11 @@ public class QueueProcessor { isProcessing = true; final T item = myQueue.removeFirst(); final Runnable runnable = () -> { - if (myDeathCondition.value(null)) return; - runSafely(() -> myProcessor.consume(item, myContinuationContext)); + if (myDeathCondition.value(null)) { + finishProcessing(false); + return; + } + runSafely(() -> myProcessor.consume(item, () -> finishProcessing(true))); }; final Application application = ApplicationManager.getApplication(); if (myThreadToUse == ThreadToUse.AWT) { diff --git a/platform/vcs-tests/testSrc/com/intellij/openapi/vcs/changes/committed/BackgroundTaskQueueTest.java b/platform/platform-tests/testSrc/com/intellij/util/concurrency/BackgroundTaskQueueTest.java similarity index 99% rename from platform/vcs-tests/testSrc/com/intellij/openapi/vcs/changes/committed/BackgroundTaskQueueTest.java rename to platform/platform-tests/testSrc/com/intellij/util/concurrency/BackgroundTaskQueueTest.java index a2fbeae9d4d4..b66b50240a11 100644 --- a/platform/vcs-tests/testSrc/com/intellij/openapi/vcs/changes/committed/BackgroundTaskQueueTest.java +++ b/platform/platform-tests/testSrc/com/intellij/util/concurrency/BackgroundTaskQueueTest.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package com.intellij.openapi.vcs.changes.committed; +package com.intellij.util.concurrency; import com.intellij.openapi.progress.BackgroundTaskQueue; import com.intellij.openapi.progress.ProcessCanceledException; diff --git a/platform/platform-tests/testSrc/com/intellij/util/concurrency/QueueProcessorTest.kt b/platform/platform-tests/testSrc/com/intellij/util/concurrency/QueueProcessorTest.kt new file mode 100644 index 000000000000..4d31434f9183 --- /dev/null +++ b/platform/platform-tests/testSrc/com/intellij/util/concurrency/QueueProcessorTest.kt @@ -0,0 +1,33 @@ +/* + * Copyright 2000-2017 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. + */ +package com.intellij.util.concurrency + +import com.intellij.testFramework.PlatformTestCase + +class QueueProcessorTest : PlatformTestCase() { + fun `test waiting for returns on finish condition`() { + var stop = false; + val semaphore = Semaphore(0) + val processor = QueueProcessor({ semaphore.down() }, { stop }) + + processor.add(1) + stop = true; + semaphore.up() + + assertTrue(processor.waitFor(1000)); + processor.waitFor(); // just in case let's check this method as well - hopefully, it won't hang since waitFor(timeout) works + } +} \ No newline at end of file