From 3fc24c1193e2a3abb91fac26b7134d18eb3410eb Mon Sep 17 00:00:00 2001 From: Eldar Abusalimov Date: Tue, 29 May 2018 13:31:36 +0300 Subject: [PATCH] util: Run QP.wrappingProcessor() continuation inside `finally` block --- .../util/concurrency/QueueProcessor.java | 10 ++-- .../util/concurrency/QueueProcessorTest.kt | 53 ++++++++++++++++++- 2 files changed, 58 insertions(+), 5 deletions(-) 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 f9c3e4a870ab..95f880bdd6de 100644 --- a/platform/platform-api/src/com/intellij/util/concurrency/QueueProcessor.java +++ b/platform/platform-api/src/com/intellij/util/concurrency/QueueProcessor.java @@ -73,9 +73,13 @@ public class QueueProcessor { @NotNull private static PairConsumer wrappingProcessor(@NotNull final Consumer processor) { - return (item, runnable) -> { - runSafely(() -> processor.consume(item)); - runnable.run(); + return (item, continuation) -> { + try { + runSafely(() -> processor.consume(item)); + } + finally { + continuation.run(); + } }; } diff --git a/platform/platform-tests/testSrc/com/intellij/util/concurrency/QueueProcessorTest.kt b/platform/platform-tests/testSrc/com/intellij/util/concurrency/QueueProcessorTest.kt index 4d31434f9183..54ac8991be5d 100644 --- a/platform/platform-tests/testSrc/com/intellij/util/concurrency/QueueProcessorTest.kt +++ b/platform/platform-tests/testSrc/com/intellij/util/concurrency/QueueProcessorTest.kt @@ -15,9 +15,17 @@ */ package com.intellij.util.concurrency +import com.intellij.execution.ExecutionException +import com.intellij.openapi.progress.ProcessCanceledException +import com.intellij.testFramework.LoggedErrorProcessor import com.intellij.testFramework.PlatformTestCase +import java.util.concurrent.LinkedBlockingQueue +import java.util.concurrent.TimeUnit + +private const val TIMEOUT_MS = 1000L class QueueProcessorTest : PlatformTestCase() { + fun `test waiting for returns on finish condition`() { var stop = false; val semaphore = Semaphore(0) @@ -26,8 +34,49 @@ class QueueProcessorTest : PlatformTestCase() { processor.add(1) stop = true; semaphore.up() - - assertTrue(processor.waitFor(1000)); + + assertTrue(processor.waitFor(TIMEOUT_MS)); processor.waitFor(); // just in case let's check this method as well - hopefully, it won't hang since waitFor(timeout) works } + + fun `test works fine after thrown exception`() { + LoggedErrorProcessor.getInstance().disableStderrDumping(testRootDisposable) + + val resultQueue = LinkedBlockingQueue() + val queueProcessor = QueueProcessor<() -> Any> { + try { + resultQueue.add(it()) + } + catch (e: Throwable) { + resultQueue.add(e) + throw e + } + } + + fun check(expectedResult: Any, item: () -> Any) { + queueProcessor.add(item) + assertEquals(expectedResult, resultQueue.poll(TIMEOUT_MS, TimeUnit.MILLISECONDS)) + assertEmpty(resultQueue) + } + fun check(expectedResult: Number) = check(expectedResult) { expectedResult } + fun check(expectedException: Throwable) = check(expectedException) { + throw expectedException.also { + it.addSuppressed(Throwable()) + } + } + + check(1) + check(Throwable()) + check(2) + check(Error()) + check(3) + check(RuntimeException()) + check(4) + check(ProcessCanceledException()) // this used to make the remaining queue elements stuck + check(5) + check(Exception()) + check(6) + check(ExecutionException("EE")) + check(7) + } } \ No newline at end of file