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
This commit is contained in:
Vladimir Krivosheev
2023-06-12 21:31:56 +00:00
committed by intellij-monorepo-bot
parent fff71c0b41
commit d8347647f0
4 changed files with 173 additions and 184 deletions
@@ -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)
@@ -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) }
}
@@ -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<PerformanceWatcherImpl>()
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<PerformanceListener>("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<FreezeCheckerTask?>(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<RegistryManager>()
.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<MySamplingTask?>()
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<StackTraceElement>? = 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<StackTraceElement>? = 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)
@@ -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()