diff --git a/java/idea-ui/src/com/intellij/ide/projectWizard/generators/SdkPreIndexingService.kt b/java/idea-ui/src/com/intellij/ide/projectWizard/generators/SdkPreIndexingService.kt index fb8884e924f4..cedeb0fa7d61 100644 --- a/java/idea-ui/src/com/intellij/ide/projectWizard/generators/SdkPreIndexingService.kt +++ b/java/idea-ui/src/com/intellij/ide/projectWizard/generators/SdkPreIndexingService.kt @@ -84,13 +84,11 @@ internal class SdkPreIndexingService: Disposable { allIterators.addAll(createIterators(sdk)) val queue = PerProjectIndexingQueue(project) - queue.getSink(0).use { sink -> - for (iterator in allIterators) { - iterator.iterateFiles(project, ContentIterator { fileOrDir: VirtualFile -> - sink.addFile(fileOrDir) - true - }, VirtualFileFilter.ALL) - } + for (iterator in allIterators) { + iterator.iterateFiles(project, ContentIterator { fileOrDir: VirtualFile -> + queue.addFile(fileOrDir, 0) + true + }, VirtualFileFilter.ALL) } return queue } diff --git a/platform/lang-impl/src/com/intellij/util/indexing/FileBasedIndexTumbler.kt b/platform/lang-impl/src/com/intellij/util/indexing/FileBasedIndexTumbler.kt index c3cd005088a1..bacacf8d8e1a 100644 --- a/platform/lang-impl/src/com/intellij/util/indexing/FileBasedIndexTumbler.kt +++ b/platform/lang-impl/src/com/intellij/util/indexing/FileBasedIndexTumbler.kt @@ -62,7 +62,6 @@ class FileBasedIndexTumbler(private val reason: @NonNls String) { scannerExecutor.cancelAllTasksAndWait() val perProjectIndexingQueue = project.getService(PerProjectIndexingQueue::class.java) - perProjectIndexingQueue.cancelAllTasksAndWait() perProjectIndexingQueue.clear() val dumbService = DumbService.getInstance(project) @@ -105,7 +104,6 @@ class FileBasedIndexTumbler(private val reason: @NonNls String) { } for (project in ProjectUtil.getOpenProjects()) { UnindexedFilesScannerExecutor.getInstance(project).resumeQueue() - project.getService(PerProjectIndexingQueue::class.java).resumeQueue() FileBasedIndexInfrastructureExtension.attachAllExtensionsData(project) } dumbModeSemaphore.up() diff --git a/platform/lang-impl/src/com/intellij/util/indexing/PerProjectIndexingQueue.kt b/platform/lang-impl/src/com/intellij/util/indexing/PerProjectIndexingQueue.kt index f95992ea67e6..6aa261c3b297 100644 --- a/platform/lang-impl/src/com/intellij/util/indexing/PerProjectIndexingQueue.kt +++ b/platform/lang-impl/src/com/intellij/util/indexing/PerProjectIndexingQueue.kt @@ -32,78 +32,9 @@ import java.util.concurrent.locks.ReadWriteLock import java.util.concurrent.locks.ReentrantReadWriteLock import kotlin.concurrent.withLock -private class PerProviderSinkFactory() { - private val activeSinksCount: AtomicInteger = AtomicInteger() - private val cancelActiveSinks: AtomicBoolean = AtomicBoolean() - - inner class PerProviderSinkImpl(private val doAddFile: (VirtualFile) -> Unit) : PerProjectIndexingQueue.PerProviderSink { - private var closed = false - - init { - activeSinksCount.incrementAndGet() - } - - override fun addFile(file: VirtualFile) { - LOG.assertTrue(!closed, "Should not invoke 'addFile' after 'close'") - if (cancelActiveSinks.get()) { - ProgressManager.getGlobalProgressIndicator()?.cancel() - ProgressManager.checkCanceled() - LOG.error("Could not cancel file addition") - } - - doAddFile(file) - } - - override fun close() { - if (!closed) { - closed = true - activeSinksCount.decrementAndGet() - } - } - } - - fun newSink(addFile: (VirtualFile) -> Unit): PerProjectIndexingQueue.PerProviderSink { - if (cancelActiveSinks.get()) { - ProgressManager.getGlobalProgressIndicator()?.cancel() - ProgressManager.checkCanceled() - LOG.error("Could not cancel sink creation") - } - - return PerProviderSinkImpl(addFile) - } - - fun cancelAllProducersAndWait() { - cancelActiveSinks.set(true) - ProgressIndicatorUtils.awaitWithCheckCanceled { - PingProgress.interactWithEdtProgress() - LockSupport.parkNanos(50_000_000) - activeSinksCount.get() == 0 - } - } - - fun resumeProducers() { - cancelActiveSinks.set(false) - } - - companion object { - private val LOG = logger() - } -} - @Internal @Service(Service.Level.PROJECT) class PerProjectIndexingQueue(private val project: Project) { - /** - * Not thread safe. These classes are cheap to construct and use - don't share instances. - *

- * Always use try-with-resources when creating instances of this interface, otherwise [cancelAllTasksAndWait] may never end waiting - */ - interface PerProviderSink : AutoCloseable { - fun addFile(file: VirtualFile) - override fun close() - } - - private val sinkFactory = PerProviderSinkFactory() @Internal class QueuedFiles { @@ -233,34 +164,14 @@ class PerProjectIndexingQueue(private val project: Project) { * Creates new instance of **thread-unsafe** [PerProviderSink] * Will throw [ProcessCanceledException] if the queue is suspended via [cancelAllTasksAndWait] */ - fun getSink(scanningId: Long): PerProviderSink { - return sinkFactory.newSink { vFile -> - // readLock here is to make sure that queuedFiles does not change during the operation - queuedFilesLock.readLock().withLock { - // .value for each file, because we want to put files into a new queue after getAndResetQueuedFiles invocation - queuedFiles.value.addFile(vFile, scanningId) - } + fun addFile(vFile: VirtualFile, scanningId: Long) { + // readLock here is to make sure that queuedFiles does not change during the operation + queuedFilesLock.readLock().withLock { + // .value for each file, because we want to put files into a new queue after getAndResetQueuedFiles invocation + queuedFiles.value.addFile(vFile, scanningId) } } - /** - * Cancels all the created [PerProviderSink] and waits until all the Sinks are finished (invoke [PerProviderSink.commit()]). - * New invocations of [PerProjectIndexingQueue.getSink()] will throw [ProcessCanceledException]. - * Use [resumeQueue] to resume the queue. - * Does nothing if the queue is already suspended. - */ - fun cancelAllTasksAndWait() { - sinkFactory.cancelAllProducersAndWait() - } - - /** - * Resumes the queue after [cancelAllTasksAndWait] invocation. - * Does nothing if the queue is already resumed. - */ - fun resumeQueue() { - sinkFactory.resumeProducers() - } - @OptIn(ExperimentalCoroutinesApi::class) fun estimatedFilesCount(): Flow = queuedFiles.flatMapLatest { it.estimatedFilesCount } diff --git a/platform/lang-impl/src/com/intellij/util/indexing/UnindexedFilesScanner.kt b/platform/lang-impl/src/com/intellij/util/indexing/UnindexedFilesScanner.kt index bbbf74233f6e..c76d5f3bdccb 100644 --- a/platform/lang-impl/src/com/intellij/util/indexing/UnindexedFilesScanner.kt +++ b/platform/lang-impl/src/com/intellij/util/indexing/UnindexedFilesScanner.kt @@ -499,49 +499,47 @@ class UnindexedFilesScanner ( sharedExplanationLogger: IndexingReasonExplanationLogger, files: ArrayDeque, ) { - project.getService(PerProjectIndexingQueue::class.java) - .getSink(scanningHistory.scanningSessionId).use { perProviderSink -> - scanningStatistics.startFileChecking() - try { - readAction { - val finder = - if (ourTestMode == TestMode.PUSHING) null - else UnindexedFilesFinder(project, sharedExplanationLogger, forceReindexingTrigger, - scanningRequest, filterHandler) - val pushingUtil = PushingUtil(project, provider) - if (!pushingUtil.mayBeUsed()) { - LOG.warn("Iterator based on $provider can't be used.") - return@readAction - } - while (files.isNotEmpty()) { - val file = files.removeFirst() - try { - if (file.isValid) { - pushingUtil.applyPushers(file) - val status = finder?.getFileStatus(file) - if (status != null) { - if (status.shouldIndex && ourTestMode == null) { - perProviderSink.addFile(file) - } - scanningStatistics.addStatus(file, status, project) - } + val indexingQueue = project.getService(PerProjectIndexingQueue::class.java) + scanningStatistics.startFileChecking() + try { + readAction { + val finder = + if (ourTestMode == TestMode.PUSHING) null + else UnindexedFilesFinder(project, sharedExplanationLogger, forceReindexingTrigger, + scanningRequest, filterHandler) + val pushingUtil = PushingUtil(project, provider) + if (!pushingUtil.mayBeUsed()) { + LOG.warn("Iterator based on $provider can't be used.") + return@readAction + } + while (files.isNotEmpty()) { + val file = files.removeFirst() + try { + if (file.isValid) { + pushingUtil.applyPushers(file) + val status = finder?.getFileStatus(file) + if (status != null) { + if (status.shouldIndex && ourTestMode == null) { + indexingQueue.addFile(file, scanningHistory.scanningSessionId) } - } - catch (e: ProcessCanceledException) { - files.addFirst(file) - throw e - } - catch (e: Exception) { - LOG.error("Error while scanning ${file.presentableUrl}\n" + - "To reindex this file IDE has to be restarted", e); + scanningStatistics.addStatus(file, status, project) } } } - } - finally { - scanningStatistics.tryFinishFilesChecking() + catch (e: ProcessCanceledException) { + files.addFirst(file) + throw e + } + catch (e: Exception) { + LOG.error("Error while scanning ${file.presentableUrl}\n" + + "To reindex this file IDE has to be restarted", e); + } } } + } + finally { + scanningStatistics.tryFinishFilesChecking() + } } private suspend fun getFilesToScan( diff --git a/platform/lang-impl/testSources/com/intellij/util/indexing/PerProviderSinkTest.kt b/platform/lang-impl/testSources/com/intellij/util/indexing/PerProviderSinkTest.kt index 8daa8c4727c7..4b8337696bf0 100644 --- a/platform/lang-impl/testSources/com/intellij/util/indexing/PerProviderSinkTest.kt +++ b/platform/lang-impl/testSources/com/intellij/util/indexing/PerProviderSinkTest.kt @@ -1,18 +1,18 @@ // Copyright 2000-2023 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license. package com.intellij.util.indexing -import com.intellij.openapi.progress.* import com.intellij.openapi.vfs.VirtualFile import com.intellij.openapi.vfs.VirtualFileWithId import com.intellij.testFramework.LightPlatformTestCase import com.intellij.testFramework.LightVirtualFile import com.intellij.util.indexing.events.FileIndexingRequest import junit.framework.TestCase -import java.util.concurrent.Phaser -import java.util.concurrent.TimeUnit +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.launch +import kotlinx.coroutines.runBlocking +import kotlinx.coroutines.sync.Semaphore import java.util.concurrent.atomic.AtomicBoolean import java.util.concurrent.atomic.AtomicInteger -import kotlin.random.Random class PerProviderSinkTest : LightPlatformTestCase() { private val DEFAULT_SCANNING_ID = 0L @@ -33,156 +33,51 @@ class PerProviderSinkTest : LightPlatformTestCase() { return LightVirtualFileWithId(name, idCounter.incrementAndGet()) } - fun testNoAddClose() { - queue.getSink(DEFAULT_SCANNING_ID).close() - val (filesInQueue, _) = getAndResetQueuedFiles(queue) - TestCase.assertEquals(0, filesInQueue.size) - } - - fun testAddClose() { - queue.getSink(DEFAULT_SCANNING_ID).use { sink -> - sink.addFile(createFile("f1")) - } - val (filesInQueue) = getAndResetQueuedFiles(queue) + fun testAddFile() { + queue.addFile(createFile("f1"), DEFAULT_SCANNING_ID) + var filesInQueue = getAndResetQueuedFiles(queue).first TestCase.assertEquals(1, filesInQueue.size) - } - fun testAddCloseClose() { - queue.getSink(DEFAULT_SCANNING_ID).use { sink -> - sink.addFile(createFile("f1")) - sink.close() - sink.close() - } - val (filesInQueue, _) = getAndResetQueuedFiles(queue) - TestCase.assertEquals(1, filesInQueue.size) - } - - fun testAddCloseTwoSinks() { - queue.getSink(DEFAULT_SCANNING_ID).use { sink -> - sink.addFile(createFile("f1")) - } - - queue.getSink(DEFAULT_SCANNING_ID).use { sink -> - sink.addFile(createFile("f2")) - } - val (filesInQueue, _) = getAndResetQueuedFiles(queue) + queue.addFile(createFile("f2"), DEFAULT_SCANNING_ID) + queue.addFile(createFile("f3"), DEFAULT_SCANNING_ID) + filesInQueue = getAndResetQueuedFiles(queue).first TestCase.assertEquals(2, filesInQueue.size) } - fun testCancelAllTasksAndWait() { - val sinkRunning = AtomicBoolean(false) - val phaser = Phaser(2) - val task = object : Task.Backgroundable(project, "Test task", true) { - override fun run(indicator: ProgressIndicator) { - assertFalse(indicator.isCanceled) - val sink = queue.getSink(DEFAULT_SCANNING_ID) - try { - phaser.awaitAdvanceInterruptibly(phaser.arrive(), 5, TimeUnit.SECONDS) // p1 - sinkRunning.set(true) - phaser.awaitAdvanceInterruptibly(phaser.arrive(), 5, TimeUnit.SECONDS) // p2 - Thread.sleep(500) // give a chance for cancelAllTasksAndWait to take effect - sink.addFile(createFile("f1")) - } - finally { - sinkRunning.set(false) - sink.close() - assertTrue(indicator.isCanceled) // sink.addFile should cancel the progress - } - } - } - - ProgressManager.getInstance().runProcessWithProgressAsynchronously(task, EmptyProgressIndicator()) - - assertFalse(sinkRunning.get()) - phaser.awaitAdvanceInterruptibly(phaser.arrive(), 5, TimeUnit.SECONDS) // p1 - phaser.awaitAdvanceInterruptibly(phaser.arrive(), 5, TimeUnit.SECONDS) // p2 - assertTrue(sinkRunning.get()) - queue.cancelAllTasksAndWait() - assertFalse(sinkRunning.get()) - } - - fun testNonCancelableSection() { - val nonCancelableSectionCompeteNormally = AtomicBoolean(false) - val phaser = Phaser(2) - val task = object : Task.Backgroundable(project, "Test task", true) { - override fun run(indicator: ProgressIndicator) { - try { - assertNotNull(ProgressManager.getGlobalProgressIndicator()) - - // There should be no PCE in non-cancelable sections - ProgressManager.getInstance().executeNonCancelableSection { - try { - assertFalse(indicator.isCanceled) - queue.getSink(DEFAULT_SCANNING_ID).use { sink -> - sink.addFile(createFile("f1")) - } - } - catch (_: ProcessCanceledException) { - fail("Should not throw PCE in non-cancellable section") - } - catch (t: Throwable) { - assertEquals("Could not cancel sink creation", t.message) - } - } - nonCancelableSectionCompeteNormally.set(true) - } - finally { - phaser.arriveAndDeregister() // p1 - } - } - } - - // cancel and wait. This will return immediately because no sinks are connected yet. But new sinks should be rejected. - queue.cancelAllTasksAndWait() - - ProgressManager.getInstance().runProcessWithProgressAsynchronously(task, EmptyProgressIndicator()) - - phaser.awaitAdvanceInterruptibly(phaser.arrive(), 5, TimeUnit.SECONDS) // p1 - assertTrue(nonCancelableSectionCompeteNormally.get()) - } - - fun testManySinksManyProvidersStress() { - val maxBatchSize = 100 + fun testStress() = runBlocking { val threadsCount = 30 val queueFlushCount = 50 val threadsCompleted = AtomicInteger() val filesSubmitted = AtomicInteger() + val producersRunning = AtomicBoolean(true) + val semaphore = Semaphore(threadsCount, threadsCount) - class SimpleRandomizedProducer(private val producerName: String) : Runnable { - override fun run() { - var batch = 0 - while (true) { - batch++ - val batchSize = Random.nextInt(maxBatchSize) - queue.getSink(DEFAULT_SCANNING_ID).use { sink -> - for (f in 1..batchSize) { - sink.addFile(createFile("$producerName batch $batch file $f")) - filesSubmitted.incrementAndGet() - } + // producers + repeat(threadsCount) { producerNr -> + launch(Dispatchers.IO) { + try { + var fileNr = 0 + while (producersRunning.get()) { + fileNr++ + queue.addFile(createFile("$producerNr file $fileNr"), DEFAULT_SCANNING_ID) + filesSubmitted.incrementAndGet() } } + finally { + threadsCompleted.incrementAndGet() + semaphore.release() + } } } - for (i in 1..threadsCount) { - Thread { - try { - ProgressManager.getInstance().runProcess(SimpleRandomizedProducer("Producer $i"), EmptyProgressIndicator()) - } - catch (_: ProcessCanceledException) { - // do nothing. This is not a production code, ignore the PCE - } - finally { - threadsCompleted.incrementAndGet() - } - }.start() - } - var totalFilesSum = 0 + // check queue state at random times while producers are running concurrently for (i in 1..queueFlushCount) { if (i == queueFlushCount) { - queue.cancelAllTasksAndWait() + // terminate all the producers and wait + producersRunning.set(false) + repeat(threadsCount) { semaphore.acquire() } } val (filesInQueue, totalFiles) = getAndResetQueuedFiles(queue) @@ -192,6 +87,7 @@ class PerProviderSinkTest : LightPlatformTestCase() { Thread.sleep(50) } + // check final state of the queue val (filesInQueue, totalFiles) = getAndResetQueuedFiles(queue) TestCase.assertEquals(0, filesInQueue.size) TestCase.assertEquals(0, totalFiles)