From d8347647f07546eab7d2faf6e7f003adff7f5f53 Mon Sep 17 00:00:00 2001 From: Vladimir Krivosheev Date: Mon, 12 Jun 2023 17:24:33 +0200 Subject: [PATCH] Optimize PerformanceWatcher implementation Remove unnecessary dispatcher and dispose method from PerformanceWatcherImpl This commit optimizes the PerformanceWatcherImpl by removing the unnecessary limitedDispatcher and the dispose method. It also simplifies the sampleJob initialization and adjusts several coroutine launches for better effectiveness. These changes improve the code readability and maintainability without changing the underlying functionality. GitOrigin-RevId: d2ae512ed95cc3a8c3c7beb1e2109401cba6b74f --- .../intellij/diagnostic/PerformanceWatcher.kt | 3 +- .../intellij/diagnostic/IdeaFreezeReporter.kt | 5 + .../diagnostic/PerformanceWatcherImpl.kt | 329 ++++++++---------- .../com/intellij/diagnostic/SamplingTask.kt | 20 +- 4 files changed, 173 insertions(+), 184 deletions(-) diff --git a/platform/ide-core-impl/src/com/intellij/diagnostic/PerformanceWatcher.kt b/platform/ide-core-impl/src/com/intellij/diagnostic/PerformanceWatcher.kt index 8f10ba6a9d2c..5592e548048c 100644 --- a/platform/ide-core-impl/src/com/intellij/diagnostic/PerformanceWatcher.kt +++ b/platform/ide-core-impl/src/com/intellij/diagnostic/PerformanceWatcher.kt @@ -1,14 +1,13 @@ // 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.diagnostic -import com.intellij.openapi.Disposable import com.intellij.openapi.application.ApplicationManager import com.intellij.openapi.components.service import org.jetbrains.annotations.ApiStatus import org.jetbrains.annotations.NonNls import java.nio.file.Path -abstract class PerformanceWatcher : Disposable { +abstract class PerformanceWatcher { interface Snapshot { fun logResponsivenessSinceCreation(activityName: @NonNls String) diff --git a/platform/platform-impl/src/com/intellij/diagnostic/IdeaFreezeReporter.kt b/platform/platform-impl/src/com/intellij/diagnostic/IdeaFreezeReporter.kt index 5cc7caec2efa..18bd4faf9acf 100644 --- a/platform/platform-impl/src/com/intellij/diagnostic/IdeaFreezeReporter.kt +++ b/platform/platform-impl/src/com/intellij/diagnostic/IdeaFreezeReporter.kt @@ -107,6 +107,11 @@ internal class IdeaFreezeReporter : PerformanceListener { super.stop() EP_NAME.forEachExtensionSafe(FreezeProfiler::stop) } + + override suspend fun stopDumpingThreads() { + super.stopDumpingThreads() + EP_NAME.forEachExtensionSafe(FreezeProfiler::stop) + } } EP_NAME.forEachExtensionSafe { it.start(reportDir) } } diff --git a/platform/platform-impl/src/com/intellij/diagnostic/PerformanceWatcherImpl.kt b/platform/platform-impl/src/com/intellij/diagnostic/PerformanceWatcherImpl.kt index 8869180edeeb..b41ac02fe87c 100644 --- a/platform/platform-impl/src/com/intellij/diagnostic/PerformanceWatcherImpl.kt +++ b/platform/platform-impl/src/com/intellij/diagnostic/PerformanceWatcherImpl.kt @@ -8,11 +8,9 @@ import com.intellij.featureStatistics.fusCollectors.LifecycleUsageTriggerCollect import com.intellij.ide.plugins.PluginManagerCore import com.intellij.internal.statistic.utils.PluginInfo import com.intellij.internal.statistic.utils.getPluginInfoByDescriptor -import com.intellij.openapi.application.ApplicationInfo -import com.intellij.openapi.application.ApplicationManager -import com.intellij.openapi.application.EDT -import com.intellij.openapi.application.PathManager +import com.intellij.openapi.application.* import com.intellij.openapi.application.impl.ApplicationInfoImpl +import com.intellij.openapi.components.serviceAsync import com.intellij.openapi.diagnostic.Attachment import com.intellij.openapi.diagnostic.Logger import com.intellij.openapi.diagnostic.getOrLogException @@ -22,9 +20,6 @@ import com.intellij.openapi.util.SystemInfo import com.intellij.openapi.util.io.FileUtilRt import com.intellij.openapi.util.io.NioFiles import com.intellij.openapi.util.registry.RegistryManager -import com.intellij.openapi.util.registry.RegistryValue -import com.intellij.openapi.util.registry.RegistryValueListener -import com.intellij.openapi.util.registry.useRegistryManagerWhenReady import com.intellij.openapi.util.text.StringUtilRt import com.intellij.util.SystemProperties import com.intellij.util.concurrency.AppExecutorUtil @@ -32,7 +27,9 @@ import com.intellij.util.concurrency.AppScheduledExecutorService import com.intellij.util.io.basicAttributesIfExists import com.intellij.util.io.sanitizeFileName import kotlinx.coroutines.* -import kotlinx.coroutines.future.asCompletableFuture +import kotlinx.coroutines.channels.BufferOverflow +import kotlinx.coroutines.flow.MutableSharedFlow +import kotlinx.coroutines.flow.collectLatest import org.jetbrains.annotations.ApiStatus import org.jetbrains.annotations.NonNls import java.io.File @@ -51,6 +48,7 @@ import kotlin.time.toDuration private val LOG: Logger get() = logger() + private const val TOLERABLE_LATENCY = 100L private const val THREAD_DUMPS_PREFIX = "threadDumps-" private const val DURATION_FILE_NAME = ".duration" @@ -59,7 +57,6 @@ private val ideStartTime = ZonedDateTime.now() private val EP_NAME = ExtensionPointName("com.intellij.idePerformanceListener") -@OptIn(ExperimentalCoroutinesApi::class) internal class PerformanceWatcherImpl(private val coroutineScope: CoroutineScope) : PerformanceWatcher() { private val logDir = PathManager.getLogDir() @@ -72,49 +69,59 @@ internal class PerformanceWatcherImpl(private val coroutineScope: CoroutineScope @Volatile private var lastSampling = System.nanoTime() private var activeEvents = 0 - private val limitedDispatcher = Dispatchers.Default.limitedParallelism(1) - private var sampleJob: Job? = null private var currentEdtEventChecker: FreezeCheckerTask? = null private val jitWatcher = JitWatcher() private val unresponsiveIntervalLazy by lazy { RegistryManager.getInstance().get("performance.watcher.unresponsive.interval.ms") } + private val isActive: Boolean = !ApplicationManager.getApplication().isHeadlessEnvironment + + private val taskFlow = MutableSharedFlow(replay = 1, onBufferOverflow = BufferOverflow.DROP_OLDEST) + init { - if (!ApplicationManager.getApplication().isHeadlessEnvironment) { + if (isActive) { coroutineScope.launch { - useRegistryManagerWhenReady { registryManager -> - val unresponsiveInterval = unresponsiveIntervalLazy - init(unresponsiveInterval, registryManager) + asyncInit() + + val samplingIntervalMs = samplingInterval + @Suppress("KotlinConstantConditions") + if (samplingIntervalMs <= 0) { + return@launch + } + + while (true) { + delay(samplingIntervalMs) + samplePerformance(samplingIntervalMs) + } + } + + coroutineScope.launch { + taskFlow.collectLatest { task -> + if (task == null) { + return@collectLatest + } + + delay(unresponsiveInterval.toLong()) + task.edtFrozen() } } } } - private suspend fun init(unresponsiveInterval: RegistryValue, registryManager: RegistryManager) { - val cancelingListener = object : RegistryValueListener { - override fun afterValueChanged(value: RegistryValue) { - LOG.info("on UI freezes more than ${unresponsiveInterval} ms will dump threads each $dumpInterval ms for $maxDumpDuration ms max") - val samplingIntervalMs = samplingInterval - sampleJob?.cancel() - @Suppress("KotlinConstantConditions") - sampleJob = if (samplingIntervalMs <= 0) { - null - } - else { - coroutineScope.launch(limitedDispatcher) { - while (true) { - delay(samplingIntervalMs) - samplePerformance(samplingIntervalMs) - } - } - } - } + private suspend fun asyncInit() { + runCatching { + reportCrashesIfAny() + }.getOrLogException(LOG) + + withContext(Dispatchers.IO) { + cleanOldFiles(logDir, 0) } - unresponsiveInterval.addListener(cancelingListener, this) + if (ApplicationInfoImpl.getShadowInstance().isEAP) { coroutineScope.launch { - val reasonableThreadPoolSize = registryManager.get("reasonable.application.thread.pool.size") + val reasonableThreadPoolSize = ApplicationManager.getApplication().serviceAsync() + .get("reasonable.application.thread.pool.size") val service = AppExecutorUtil.getAppScheduledExecutorService() as AppScheduledExecutorService val allAvailableProcessors = Runtime.getRuntime().availableProcessors() service.setNewThreadListener { _, _ -> @@ -128,15 +135,6 @@ internal class PerformanceWatcherImpl(private val coroutineScope: CoroutineScope } } } - - runCatching { - reportCrashesIfAny() - }.getOrLogException(LOG) - - withContext(Dispatchers.IO) { - cleanOldFiles(logDir, 0) - } - cancelingListener.afterValueChanged(unresponsiveInterval) } override suspend fun processUnfinishedFreeze(consumer: suspend (Path, Int) -> Unit) { @@ -169,10 +167,6 @@ internal class PerformanceWatcherImpl(private val coroutineScope: CoroutineScope } } - override fun dispose() { - sampleJob?.cancel() - } - @Suppress("SameParameterValue") private suspend fun samplePerformance(samplingIntervalMs: Long) { val current = System.nanoTime() @@ -186,7 +180,7 @@ internal class PerformanceWatcherImpl(private val coroutineScope: CoroutineScope diffMs -= samplingIntervalMs } jitWatcher.checkJitState() - val latencyMs = withContext(Dispatchers.EDT) { + val latencyMs = withContext(Dispatchers.EDT + ModalityState.any().asContextElement()) { TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - current) } swingApdex = swingApdex.withEvent(TOLERABLE_LATENCY, latencyMs) @@ -215,21 +209,29 @@ internal class PerformanceWatcherImpl(private val coroutineScope: CoroutineScope @ApiStatus.Internal override fun edtEventStarted() { + if (!isActive) { + return + } + val start = System.nanoTime() activeEvents++ - if (sampleJob != null) { - currentEdtEventChecker?.stop() - currentEdtEventChecker = FreezeCheckerTask(start) - } + currentEdtEventChecker?.stop() + val task = FreezeCheckerTask(start) + currentEdtEventChecker = task + check(taskFlow.tryEmit(task)) } @ApiStatus.Internal override fun edtEventFinished() { - activeEvents-- - if (sampleJob != null) { - currentEdtEventChecker?.stop() - currentEdtEventChecker = if (activeEvents > 0) FreezeCheckerTask(System.nanoTime()) else null + if (!isActive) { + return } + + activeEvents-- + currentEdtEventChecker?.stop() + val task = if (activeEvents > 0) FreezeCheckerTask(System.nanoTime()) else null + currentEdtEventChecker = task + check(taskFlow.tryEmit(task)) } override fun dumpThreads(pathPrefix: String, appendMillisecondsToFileName: Boolean, stripDump: Boolean): Path? { @@ -240,9 +242,6 @@ internal class PerformanceWatcherImpl(private val coroutineScope: CoroutineScope } private fun doDumpThreads(pathPrefix: String, appendMillisecondsToFileName: Boolean, contentsPrefix: String, stripDump: Boolean): Path? { - if (sampleJob == null) { - return null - } return dumpThreads(pathPrefix = pathPrefix, appendMillisecondsToFileName = appendMillisecondsToFileName, rawDump = contentsPrefix + ThreadDumper.getThreadDumpInfo(ThreadDumper.getThreadInfos(), stripDump).rawDump) @@ -297,155 +296,87 @@ internal class PerformanceWatcherImpl(private val coroutineScope: CoroutineScope override fun clearFreezeStacktraces() { coroutineScope.launch { - currentEdtEventChecker?.stopDumping() + currentEdtEventChecker?.stopDumpingAsync() } } private inner class FreezeCheckerTask(private val taskStart: Long) { private val state = AtomicReference(CheckerState.CHECKING) - private val job: Job - private var freezeFolder: String? = null - - @Volatile - private var dumpTask: SamplingTask? = null - - init { - job = coroutineScope.launch(limitedDispatcher) { - delay(unresponsiveInterval.toLong()) - edtFrozen() - } - } - - private fun getDuration(current: Long, unit: TimeUnit): Long { - return unit.convert(current - taskStart, TimeUnit.NANOSECONDS) - } + private val dumpTask = AtomicReference() fun stop() { - job.cancel() - if (state.getAndSet(CheckerState.FINISHED) == CheckerState.FREEZE) { - val taskStop = System.nanoTime() - // stop sampling as early as possible - stopDumping() - try { - coroutineScope.launch(limitedDispatcher) { - stopDumping() - val durationMs = getDuration(taskStop, TimeUnit.MILLISECONDS) + if (state.getAndSet(CheckerState.FINISHED) != CheckerState.FREEZE) { + return + } - val freezeDir = logDir.resolve(freezeFolder!!) - for (listener in EP_NAME.extensionList) { - listener.uiFreezeFinished(durationMs, freezeDir) - } - publisher?.uiFreezeFinished(durationMs, freezeDir) + val taskStop = System.nanoTime() + val task = dumpTask.getAndSet(null) ?: return + coroutineScope.launch { + task.stopDumpingThreads() - val reportDir = postProcessReportFolder(durationMs) + val durationMs = TimeUnit.MILLISECONDS.convert(taskStop - taskStart, TimeUnit.NANOSECONDS) - for (listener in EP_NAME.extensionList) { - listener.uiFreezeRecorded(durationMs, reportDir) - } - }.asCompletableFuture().join() + val freezeFolder = task.freezeFolder + val freezeDir = logDir.resolve(freezeFolder) + for (listener in EP_NAME.extensionList) { + listener.uiFreezeFinished(durationMs, freezeDir) } - catch (e: Exception) { - LOG.warn(e) + publisher?.uiFreezeFinished(durationMs, freezeDir) + + val reportDir = postProcessReportFolder(durationMs = durationMs, task = task, dir = logDir.resolve(freezeFolder), logDir = logDir) + + for (listener in EP_NAME.extensionList) { + listener.uiFreezeRecorded(durationMs, reportDir) } } } - private suspend fun edtFrozen() { - freezeFolder = "${THREAD_DUMPS_PREFIX}freeze-${formatTime(ZonedDateTime.now())}-${buildName()}" + fun edtFrozen() { if (!state.compareAndSet(CheckerState.CHECKING, CheckerState.FREEZE)) { return } + val freezeFolder = "${THREAD_DUMPS_PREFIX}freeze-${formatTime(ZonedDateTime.now())}-${buildName()}" + //TODO always true for some reason //myFreezeDuringStartup = !LoadingState.INDEXING_FINISHED.isOccurred(); - val reportDir = logDir.resolve(freezeFolder!!) + val reportDir = logDir.resolve(freezeFolder) Files.createDirectories(reportDir) for (listener in EP_NAME.extensionList) { listener.uiFreezeStarted(reportDir, coroutineScope) } + dumpTask.set(MySamplingTask(freezeFolder = freezeFolder, taskStart = taskStart)) publisher?.uiFreezeStarted(reportDir) - - dumpTask = object : SamplingTask(dumpInterval = dumpInterval, maxDurationMs = maxDumpDuration, coroutineScope = coroutineScope) { - override suspend fun dumpedThreads(threadDump: ThreadDump) { - if (state.get() == CheckerState.FINISHED) { - stop() - return - } - - val file = dumpThreads(pathPrefix = "$freezeFolder/", appendMillisecondsToFileName = false, rawDump = threadDump.rawDump) - ?: return - try { - val duration = getDuration(System.nanoTime(), TimeUnit.SECONDS) - withContext(Dispatchers.IO) { - Files.createDirectories(file.parent) - Files.writeString(file.parent.resolve(DURATION_FILE_NAME), duration.toString()) - } - - for (listener in EP_NAME.extensionList) { - coroutineContext.ensureActive() - listener.dumpedThreads(file, threadDump) - } - coroutineContext.ensureActive() - publisher?.dumpedThreads(file, threadDump) - } - catch (e: IOException) { - LOG.info("Failed to write the duration file", e) - } - } - } } - private fun postProcessReportFolder(durationMs: Long): Path? { - val dir = logDir.resolve(freezeFolder!!) - if (!Files.exists(dir)) { - return null - } + suspend fun stopDumpingAsync() { + (dumpTask.getAndSet(null) ?: return).stopDumpingThreads() + } + } - cleanup(dir) - var reportDir = logDir.resolve("${dir.name}${getFreezePlaceSuffix()}-${TimeUnit.MILLISECONDS.toSeconds(durationMs)}sec") + inner class MySamplingTask(@JvmField val freezeFolder: String, private val taskStart: Long) + : SamplingTask(dumpInterval = dumpInterval, maxDurationMs = maxDumpDuration, coroutineScope = coroutineScope) { + override suspend fun dumpedThreads(threadDump: ThreadDump) { + val file = dumpThreads(pathPrefix = "$freezeFolder/", appendMillisecondsToFileName = false, rawDump = threadDump.rawDump) ?: return try { - Files.move(dir, reportDir) + val durationInSeconds = TimeUnit.SECONDS.convert(System.nanoTime() - taskStart, TimeUnit.NANOSECONDS) + withContext(Dispatchers.IO) { + val parent = file.parent + Files.createDirectories(parent) + Files.writeString(parent.resolve(DURATION_FILE_NAME), durationInSeconds.toString()) + } + + for (listener in EP_NAME.extensionList) { + coroutineContext.ensureActive() + listener.dumpedThreads(file, threadDump) + } + coroutineContext.ensureActive() + publisher?.dumpedThreads(file, threadDump) } catch (e: IOException) { - LOG.warn("Unable to create freeze folder $reportDir", e) - reportDir = dir + LOG.info("Failed to write the duration file", e) } - val message = "UI was frozen for ${durationMs}ms, details saved to $reportDir" - if (PluginManagerCore.isRunningFromSources()) { - LOG.info(message) - } - else { - LOG.warn(message) - } - return reportDir - } - - fun stopDumping() { - val task = dumpTask ?: return - task.stop() - dumpTask = null - } - - private fun getFreezePlaceSuffix(): String { - val task = dumpTask ?: return "" - var stacktraceCommonPart: List? = null - for (info in task.threadInfos) { - val edt = info.firstOrNull(ThreadDumper::isEDT) ?: continue - val edtStack = edt.stackTrace ?: continue - stacktraceCommonPart = if (stacktraceCommonPart == null) { - edtStack.toList() - } - else { - getStacktraceCommonPart(stacktraceCommonPart, edtStack) - } - } - - if (!stacktraceCommonPart.isNullOrEmpty()) { - val element = stacktraceCommonPart[0] - return "-${sanitizeFileName(StringUtilRt.getShortName(element.className))}.${sanitizeFileName(element.methodName)}" - } - return "" } } @@ -468,6 +399,52 @@ internal class PerformanceWatcherImpl(private val coroutineScope: CoroutineScope } } +private fun postProcessReportFolder(durationMs: Long, task: SamplingTask, dir: Path, logDir: Path): Path? { + if (!Files.exists(dir)) { + return null + } + + cleanup(dir) + var reportDir = logDir.resolve("${dir.name}${getFreezePlaceSuffix(task)}-${TimeUnit.MILLISECONDS.toSeconds(durationMs)}sec") + try { + Files.move(dir, reportDir) + } + catch (e: IOException) { + LOG.warn("Unable to create freeze folder $reportDir", e) + reportDir = dir + } + + val message = "UI was frozen for ${durationMs}ms, details saved to $reportDir" + if (PluginManagerCore.isRunningFromSources()) { + LOG.info(message) + } + else { + LOG.warn(message) + } + return reportDir +} + +private fun getFreezePlaceSuffix(task: SamplingTask): String { + var stacktraceCommonPart: List? = null + for (info in task.threadInfos) { + val edt = info.firstOrNull(ThreadDumper::isEDT) ?: continue + val edtStack = edt.stackTrace ?: continue + stacktraceCommonPart = if (stacktraceCommonPart == null) { + edtStack.toList() + } + else { + getStacktraceCommonPart(stacktraceCommonPart, edtStack) + } + } + + if (stacktraceCommonPart.isNullOrEmpty()) { + return "" + } + + val element = stacktraceCommonPart[0] + return "-${sanitizeFileName(StringUtilRt.getShortName(element.className))}.${sanitizeFileName(element.methodName)}" +} + private suspend fun reportCrashesIfAny() { val systemDir = Path.of(PathManager.getSystemPath()) val appInfoFile = systemDir.resolve(APP_INFO_FILE_NAME) diff --git a/platform/platform-impl/src/com/intellij/diagnostic/SamplingTask.kt b/platform/platform-impl/src/com/intellij/diagnostic/SamplingTask.kt index 8c171b744634..34a97ba402b6 100644 --- a/platform/platform-impl/src/com/intellij/diagnostic/SamplingTask.kt +++ b/platform/platform-impl/src/com/intellij/diagnostic/SamplingTask.kt @@ -33,25 +33,29 @@ internal open class SamplingTask(@JvmField internal val dumpInterval: Int, maxDu job = coroutineScope.launch { val delayDuration = dumpInterval.milliseconds while (true) { - dumpThreads() + dumpThreads(asyncCoroutineScope = coroutineScope) delay(delayDuration) } } } - private suspend fun dumpThreads() { + private suspend fun dumpThreads(asyncCoroutineScope: CoroutineScope) { currentTime = System.nanoTime() gcCurrentTime = currentGcTime() val infos = ThreadDumper.getThreadInfos(THREAD_MX_BEAN, false) coroutineContext.ensureActive() threadInfos = threadInfos.add(infos) - if (threadInfos.size >= maxDumps) { - stop() + if (threadInfos.size > maxDumps) { + stopDumpingThreads() + return } - coroutineContext.ensureActive() - dumpedThreads(ThreadDumper.getThreadDumpInfo(infos, true)) + asyncCoroutineScope.launch { + dumpedThreads(ThreadDumper.getThreadDumpInfo(infos, true)) + } + // don't schedule yet another dumpedThreads - wait for completion + .join() } protected open suspend fun dumpedThreads(threadDump: ThreadDump) {} @@ -63,6 +67,10 @@ internal open class SamplingTask(@JvmField internal val dumpInterval: Int, maxDu open fun stop() { job?.cancel() } + + open suspend fun stopDumpingThreads() { + job?.cancelAndJoin() + } } private val THREAD_MX_BEAN = ManagementFactory.getThreadMXBean()