IJPL-3020: simplify PerProjectIndexingQueue (per-provider sinks are not needed)

Sinks were initially introduced to reduce contention on lock that protects shared files queue. For some time already sinks work as a light-weight adapter for the shared queue (meaning that locks is used). We've never seen this lock in profiles since then.

There is a functional change however: cancelAllTasksAndWait method. It is not needed anymore: scanning should be cancelled through scanning executor (and it is cancelled through executor indeed for quite some time already).

GitOrigin-RevId: 3677aafb40e6eb47aa5b2f52001eef5cf88d0738
This commit is contained in:
Andrei.Kuznetsov
2025-01-03 14:14:03 +00:00
committed by intellij-monorepo-bot
parent 85c069ab11
commit 8cd15cf1e6
5 changed files with 76 additions and 275 deletions
@@ -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
}
@@ -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()
@@ -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<PerProviderSinkFactory>()
}
}
@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.
* <p>
* 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<Int> = queuedFiles.flatMapLatest { it.estimatedFilesCount }
@@ -499,49 +499,47 @@ class UnindexedFilesScanner (
sharedExplanationLogger: IndexingReasonExplanationLogger,
files: ArrayDeque<VirtualFile>,
) {
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(
@@ -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)