[python] PY-85791: Introduce ExecOptions.weight to limit concurrent processes.

Even though project configurators are light, we still do not want to run 2000 `hatch` concurrently.

Conda is heavy, no more than 2 of them should ever be launched to freeze user OS.

We introduce a dead-simple API (see `ExecOptions.weight`) and use it to solve aforementioned cases.

GitOrigin-RevId: da20031c74824fc36288ddb57639e07d3a9ea5a3
This commit is contained in:
Ilya.Kazakevich
2025-11-29 23:31:55 +00:00
committed by intellij-monorepo-bot
parent 2c3ef39f50
commit 40bb816185
9 changed files with 183 additions and 55 deletions
@@ -18,4 +18,18 @@
/>
</extensions>
<extensions defaultExtensionNs="com.intellij">
<registryKey defaultValue="2"
description="How many heavy processes to execute concurrently"
key="python.execService.limit.heavy" restartRequired="true"/>
<registryKey defaultValue="10"
description="How many medium processes to execute concurrently"
key="python.execService.limit.medium" restartRequired="true"/>
<registryKey defaultValue="100"
description="How many light processes to execute concurrently"
key="python.execService.limit.light" restartRequired="true"/>
</extensions>
</idea-plugin>
@@ -6,6 +6,8 @@ import com.intellij.execution.process.ProcessOutputTypes
import com.intellij.execution.target.FullPathOnTarget
import com.intellij.execution.target.TargetEnvironmentConfiguration
import com.intellij.execution.target.TargetedCommandLineBuilder
import com.intellij.openapi.application.ApplicationManager
import com.intellij.openapi.components.service
import com.intellij.openapi.util.NlsSafe
import com.intellij.platform.eel.EelApi
import com.intellij.platform.eel.getShell
@@ -33,7 +35,7 @@ import kotlin.time.Duration.Companion.minutes
/**
* Default service implementation
*/
fun ExecService(): ExecService = ExecServiceImpl
fun ExecService(): ExecService = ApplicationManager.getApplication().service<ExecServiceImpl>()
/**
@@ -219,18 +221,29 @@ open class ZeroCodeStdoutParserTransformer<T>(val stdoutParser: (String) -> Resu
}
}
/**
* Each process launched with [ExecOptions] belongs to one of these categories. The lighter proces is, the more processes system can run.
* Limits are set via Registry.
*/
enum class ConcurrentProcessWeight {
LIGHT,
MEDIUM,
HEAVY
}
/**
* @property[env] Environment variables to be applied with the process run
* @property[timeout] Process gets killed after this timeout
* @property[processDescription] optional description to be displayed to user
* @property[tty] Much like [com.intellij.platform.eel.EelExecApi.Pty]
* @property[weight] use it to limit the number of concurrent processes not to exhaust user resources, see [ConcurrentProcessWeight]
*/
data class ExecOptions(
override val env: Map<String, String> = emptyMap(),
override val processDescription: @Nls String? = null,
val timeout: Duration = 5.minutes,
override val tty: TtySize? = null,
val weight: ConcurrentProcessWeight = ConcurrentProcessWeight.LIGHT,
) : ExecOptionsBase
@@ -0,0 +1,31 @@
// Copyright 2000-2025 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license.
package com.intellij.python.community.execService.impl
import com.intellij.openapi.diagnostic.fileLogger
import com.intellij.python.community.execService.ConcurrentProcessWeight
import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.Semaphore
import kotlinx.coroutines.sync.withLock
import kotlinx.coroutines.sync.withPermit
import java.util.*
internal class ConcurrencyLimiter(private val getBucketSize: (ConcurrentProcessWeight) -> Int) {
companion object {
private val logger = fileLogger()
}
// Registry might not be accessible at the early stage, but once this object is constructed, the app is started, hence there is a registry.
private val semaphores = EnumMap<ConcurrentProcessWeight, Semaphore>(ConcurrentProcessWeight::class.java)
private val accessSemaMutex = Mutex()
suspend fun <T> exec(weight: ConcurrentProcessWeight, code: suspend () -> T): T {
val sema = accessSemaMutex.withLock { semaphores.getOrPut(weight) { Semaphore(getBucketSize(weight)) } }
return sema.withPermit {
if (logger.isTraceEnabled) {
logger.trace("For $weight perms. left: ${sema.availablePermits}")
}
code()
}
}
}
@@ -5,6 +5,7 @@ import com.intellij.openapi.application.ApplicationManager
import com.intellij.openapi.components.Service
import com.intellij.openapi.components.service
import com.intellij.openapi.diagnostic.fileLogger
import com.intellij.openapi.util.registry.Registry
import com.intellij.python.community.execService.*
import com.intellij.python.community.execService.impl.processLaunchers.LaunchRequest
import com.intellij.python.community.execService.impl.processLaunchers.ProcessLauncher
@@ -19,7 +20,15 @@ import kotlinx.coroutines.withTimeout
import org.jetbrains.annotations.Nls
internal object ExecServiceImpl : ExecService {
@Service(Service.Level.APP)
internal class ExecServiceImpl private constructor() : ExecService {
private val execLimiter = ConcurrencyLimiter {
when (it) {
ConcurrentProcessWeight.LIGHT -> Registry.intValue("python.execService.limit.light")
ConcurrentProcessWeight.MEDIUM -> Registry.intValue("python.execService.limit.medium")
ConcurrentProcessWeight.HEAVY -> Registry.intValue("python.execService.limit.heavy")
}
}
override suspend fun executeGetProcess(binary: BinaryToExec, args: Args, scopeToBind: CoroutineScope?, options: ExecGetProcessOptions): Result<Process, ExecuteGetProcessError<*>> {
val launcher = create(binary, args, options, scopeToBind).getOr { return it }
@@ -51,46 +60,57 @@ internal object ExecServiceImpl : ExecService {
return@coroutineScope it.asPyError()
}
val description = options.processDescription
?: PyExecBundle.message("py.exec.defaultName.process", (listOf(processLauncher.exeForError.toString()) + processLauncher.args).joinToString(" "))
val process = processLauncher.start().getOr {
val message = PyExecBundle.message("py.exec.start.error", description, it.error.cantExecProcessError, it.error.errNo
?: "unknown")
return@coroutineScope processLauncher.createExecError(
messageToUser = message,
errorReason = it.error
)
return@coroutineScope execLimiter.exec(options.weight) {
executeProcessGetResultImpl(binary, options, processLauncher, processInteractiveHandler)
}
val result = try {
withTimeout(options.timeout) {
val interactiveResult = processInteractiveHandler.getResultFromProcess(binary, processLauncher.args, process)
val successResult = interactiveResult.getOr { failure ->
val (output, customErrorMessage) = failure.error
val additionalMessage = customErrorMessage ?: run {
PyExecBundle.message("py.exec.exitCode.error", description, output.exitCode)
}
return@withTimeout processLauncher.createExecError(
messageToUser = additionalMessage,
errorReason = ExecErrorReason.UnexpectedProcessTermination(output),
loggedProcessId = process.loggedProcess.id,
)
}
Result.success(successResult)
}
}
catch (_: TimeoutCancellationException) {
processLauncher.killAndJoin()
processLauncher.createExecError(
messageToUser = PyExecBundle.message("py.exec.timeout.error", description, options.timeout),
errorReason = ExecErrorReason.Timeout,
loggedProcessId = process.loggedProcess.id,
)
}
return@coroutineScope result
}
}
private suspend fun <T> executeProcessGetResultImpl(
binary: BinaryToExec,
options: ExecOptions,
processLauncher: ProcessLauncher,
processInteractiveHandler: ProcessInteractiveHandler<T>,
): Result<T, ExecError> {
val description = options.processDescription
?: PyExecBundle.message("py.exec.defaultName.process", (listOf(processLauncher.exeForError.toString()) + processLauncher.args).joinToString(" "))
val process = processLauncher.start().getOr {
val message = PyExecBundle.message("py.exec.start.error", description, it.error.cantExecProcessError, it.error.errNo
?: "unknown")
return processLauncher.createExecError(
messageToUser = message,
errorReason = it.error
)
}
val result = try {
withTimeout(options.timeout) {
val interactiveResult = processInteractiveHandler.getResultFromProcess(binary, processLauncher.args, process)
val successResult = interactiveResult.getOr { failure ->
val (output, customErrorMessage) = failure.error
val additionalMessage = customErrorMessage ?: run {
PyExecBundle.message("py.exec.exitCode.error", description, output.exitCode)
}
return@withTimeout processLauncher.createExecError(
messageToUser = additionalMessage,
errorReason = ExecErrorReason.UnexpectedProcessTermination(output),
loggedProcessId = process.loggedProcess.id,
)
}
Result.success(successResult)
}
}
catch (_: TimeoutCancellationException) {
processLauncher.killAndJoin()
processLauncher.createExecError(
messageToUser = PyExecBundle.message("py.exec.timeout.error", description, options.timeout),
errorReason = ExecErrorReason.Timeout,
loggedProcessId = process.loggedProcess.id,
)
}
return result
}
}
private fun <T : ExecErrorReason> ProcessLauncher.createExecError(
@@ -0,0 +1,60 @@
// Copyright 2000-2025 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license.
package com.intellij.python.junit5Tests.unit
import com.intellij.execution.Platform
import com.intellij.openapi.util.registry.Registry
import com.intellij.platform.eel.getShell
import com.intellij.platform.eel.provider.asNioPath
import com.intellij.platform.eel.provider.localEel
import com.intellij.platform.eel.provider.systemOs
import com.intellij.python.community.execService.*
import com.intellij.testFramework.common.timeoutRunBlocking
import com.intellij.testFramework.junit5.TestApplication
import kotlinx.coroutines.*
import kotlinx.coroutines.sync.Semaphore
import org.junit.jupiter.api.Assertions
import org.junit.jupiter.api.Test
import org.junit.jupiter.api.io.TempDir
import java.nio.file.Path
import kotlin.io.path.listDirectoryEntries
import kotlin.time.Duration.Companion.milliseconds
import kotlin.time.Duration.Companion.minutes
@TestApplication
class ExecLimiterTest {
@Test
fun testLimit(@TempDir dir: Path): Unit = timeoutRunBlocking(1.minutes) {
val (shell, arg) = localEel.exec.getShell()
val sleepCmd = when (localEel.systemOs().platform) {
Platform.WINDOWS -> "pause"
// sleep(1) can't be used as it creates a separate process which can't be killed as ijent doesn't kill the process group
// cat, from the other hand, stops as soon as its stdin gets closed, so killing sh is enough
Platform.UNIX -> "cat"
}
val execService = ExecService()
val weight = ConcurrentProcessWeight.HEAVY
val limit = Registry.intValue("python.execService.limit.heavy")
val workersCount = limit * 10
val lock = Semaphore(permits = workersCount, acquiredPermits = workersCount)
repeat(workersCount) { n ->
launch(Dispatchers.Unconfined, start = CoroutineStart.UNDISPATCHED) {
lock.release()
val file = dir.resolve("$n.txt")
execService.execGetStdout(
binary = shell.asNioPath(),
args = Args(arg, "echo 1 > $file && $sleepCmd"),
options = ExecOptions(weight = weight)).orThrow()
}
}
repeat(workersCount) {
lock.acquire()
}
delay(500.milliseconds)
Assertions.assertEquals(limit, dir.listDirectoryEntries().size, "No more than $limit processes must run at the same time")
coroutineContext.job.cancelChildren()
}
}
@@ -17,10 +17,6 @@
<extensions defaultExtensionNs="com.intellij.platform">
<rpc.backend.remoteApiProvider implementation="com.intellij.python.sdkConfigurator.backend.impl.rpcBridge.ApiProvider"/>
</extensions>
<extensions defaultExtensionNs="com.intellij">
<registryKey defaultValue="1" description="How many SDK lookup processes execute in parallel"
key="intellij.python.sdkConfigurator.backend.sdk.parallel"/>
</extensions>
<actions>
<action class="com.intellij.python.sdkConfigurator.backend.impl.platformBridge.ConfigureSDKAction" id="ConfigureSDKAction"/>
@@ -6,10 +6,8 @@ import com.intellij.openapi.module.Module
import com.intellij.openapi.project.Project
import com.intellij.openapi.project.guessModuleDir
import com.intellij.openapi.project.modules
import com.intellij.openapi.roots.ModuleRootManager
import com.intellij.openapi.roots.ModuleRootModificationUtil
import com.intellij.openapi.util.Key
import com.intellij.openapi.util.registry.Registry
import com.intellij.openapi.util.removeUserData
import com.intellij.platform.ide.progress.withBackgroundProgress
import com.intellij.python.common.tools.ToolId
@@ -32,8 +30,6 @@ import kotlinx.collections.immutable.toPersistentList
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.async
import kotlinx.coroutines.awaitAll
import kotlinx.coroutines.sync.Semaphore
import kotlinx.coroutines.sync.withPermit
import kotlinx.coroutines.withContext
/**
@@ -107,14 +103,11 @@ internal class ModulesSdkConfigurator private constructor(
private suspend fun getModulesWithoutSDKCreateInfo(project: Project): Map<ModuleName, ModuleCreateInfo> = withBackgroundProgress(project, PySdkConfiguratorBundle.message("intellij.python.sdk.looking")) {
val tools = PyProjectSdkConfigurationExtension.createMap()
val limit = Semaphore(permits = Registry.intValue("intellij.python.sdkConfigurator.backend.sdk.parallel"))
val now = System.currentTimeMillis()
val resultDef = project.modules.filter { PythonSdkUtil.findPythonSdk(it) == null }.map { module ->
limit.withPermit {
async {
val moduleInfo = getModuleInfo(module, tools) ?: return@async null
Pair(module, moduleInfo)
}
async {
val moduleInfo = getModuleInfo(module, tools) ?: return@async null
Pair(module, moduleInfo)
}
}
val result = resultDef.awaitAll().filterNotNull()
@@ -26,8 +26,9 @@ suspend fun <T> runExecutableWithProgress(
env: Map<String, String> = emptyMap(),
vararg args: String,
transformer: ProcessOutputTransformer<T>,
processWeight: ConcurrentProcessWeight = ConcurrentProcessWeight.LIGHT
): PyResult<T> {
val execOptions = ExecOptions(timeout = timeout, env = env)
val execOptions = ExecOptions(timeout = timeout, env = env, weight = processWeight)
val errorHandlerTransformer: ProcessOutputTransformer<T> = { output ->
when {
@@ -133,7 +133,7 @@ object CondaExecutor {
): PyResult<T> {
val envs = getFixedEnvs(binaryToExec).getOr { return it }
val runArgs = prepareCondaRunArgs(args, emptyList(), condaEnvIdentity).toTypedArray()
return runExecutableWithProgress(binaryToExec, timeout, env = envs, *runArgs, transformer = transformer)
return runExecutableWithProgress(binaryToExec, timeout, env = envs, *runArgs, transformer = transformer, processWeight = ConcurrentProcessWeight.HEAVY)
}
private fun getFixedEnvs(binaryToExec: BinaryToExec): PyResult<Map<String, String>> {