From 8750ef2fe491536129c08e150950992d71e91b92 Mon Sep 17 00:00:00 2001 From: Sergey Simonchik Date: Tue, 28 Nov 2023 23:21:05 +0100 Subject: [PATCH] terminal: make sure generators and user commands don't interfere with each other (IDEA-337968) GitOrigin-RevId: dbcda6df0dd0e5dccaaf217e195aecbede173a41 --- .../exp/ShellCommandExecutionManager.kt | 186 ++++++++++++++++++ .../terminal/exp/ShellCommandManager.kt | 15 +- .../plugins/terminal/exp/TerminalSession.kt | 12 +- .../completion/IJShellRuntimeDataProvider.kt | 28 +-- .../terminal/block/BlockTerminalTest.kt | 34 +++- 5 files changed, 236 insertions(+), 39 deletions(-) create mode 100644 plugins/terminal/src/org/jetbrains/plugins/terminal/exp/ShellCommandExecutionManager.kt diff --git a/plugins/terminal/src/org/jetbrains/plugins/terminal/exp/ShellCommandExecutionManager.kt b/plugins/terminal/src/org/jetbrains/plugins/terminal/exp/ShellCommandExecutionManager.kt new file mode 100644 index 000000000000..9e540394d62c --- /dev/null +++ b/plugins/terminal/src/org/jetbrains/plugins/terminal/exp/ShellCommandExecutionManager.kt @@ -0,0 +1,186 @@ +// Copyright 2000-2023 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license. +package org.jetbrains.plugins.terminal.exp + +import com.intellij.util.containers.nullize +import com.intellij.util.execution.ParametersListUtil +import kotlinx.coroutines.CompletableDeferred +import org.jetbrains.annotations.TestOnly +import org.jetbrains.plugins.terminal.TerminalUtil +import org.jetbrains.plugins.terminal.exp.ShellCommandManager.Companion.LOG +import java.util.* +import java.util.concurrent.CancellationException +import java.util.concurrent.atomic.AtomicInteger +import kotlin.collections.ArrayList + +internal class ShellCommandExecutionManager(private val session: TerminalSession, commandManager: ShellCommandManager) { + + private val lock: Lock = Lock() + + // these fields are guarded by `lock` + private val scheduledGenerators: Queue = LinkedList() + private var runningGenerator: Generator? = null + private val scheduledCommands: Queue = LinkedList() + private var isCommandRunning: Boolean = false + + @Volatile + private var generatorCommandSent: CompletableDeferred = CompletableDeferred() + + init { + commandManager.addListener(object : ShellCommandListener { + override fun commandFinished(command: String?, exitCode: Int, duration: Long?) { + lock.withLock { withoutLock -> + if (!isCommandRunning) { + LOG.warn("Received command_finished event, but command wasn't started") + } + isCommandRunning = false + if (runningGenerator != null) { + val runningGeneratorLocal = runningGenerator!! + runningGenerator = null + withoutLock { + val msg = "Unexpectedly running $runningGeneratorLocal when command_finished event received" + LOG.warn(msg) + runningGeneratorLocal.deferred.completeExceptionally(IllegalStateException(msg)) + } + } + scheduledGenerators.drainToList().nullize()?.let { cancelledGenerators -> + LOG.warn("Unexpected scheduled generators $cancelledGenerators when command_finished event received") + withoutLock { + cancelledGenerators.forEach { + it.deferred.cancel(CancellationException( + "Unexpected scheduled generators when command_finished event received")) + } + } + } + } + processQueueIfReady() + } + + override fun generatorFinished(requestId: Int, result: String) { + lock.withLock { withoutLock -> + if (runningGenerator == null) { + LOG.warn("Received generator_finished event (request_id=${requestId}), but no running generator") + } + else { + val runningGeneratorLocal = runningGenerator!! + runningGenerator = null + withoutLock { + if (requestId == runningGeneratorLocal.requestId) { + runningGeneratorLocal.deferred.complete(result) + } + else { + val msg = "Received generator_finished event (request_id=${requestId}), but $runningGeneratorLocal was expected" + LOG.warn(msg) + runningGeneratorLocal.deferred.completeExceptionally(IllegalStateException(msg)) + } + } + } + } + processQueueIfReady() + } + }, session) + } + + fun sendCommandToExecute(shellCommand: String) { + lock.withLock { + if (isCommandRunning) { + LOG.warn("Command '$shellCommand' execution is postponed until currently running command is finished") + } + scheduledCommands.offer(shellCommand) + } + processQueueIfReady() + } + + fun runGeneratorAsync(generatorName: String, generatorParameters: List): CompletableDeferred { + val generator = Generator(generatorName, generatorParameters) + lock.withLock { withoutLock -> + if (isCommandRunning) { + withoutLock { + generator.deferred.completeExceptionally(IllegalStateException( + "Generator shouldn't be scheduled when command is running" + )) + } + } + scheduledGenerators.offer(generator) + } + processQueueIfReady() + return generator.deferred + } + + @TestOnly + suspend fun awaitGeneratorCommandSent() { + generatorCommandSent.await() + } + + // should be called without `lock` + private fun processQueueIfReady() { + lock.withLock { withoutLock -> + if (runningGenerator == null && !isCommandRunning) { + scheduledCommands.poll()?.let { command -> + // cancel previously scheduled generators, because user command is already ready + scheduledGenerators.drainToList().nullize()?.let { cancelledGenerators -> + withoutLock { + cancelledGenerators.forEach { it.deferred.cancel(CancellationException("Generator cancelled because of executing command")) } + } + } + isCommandRunning = true + doSendCommandToExecute(command) + return@withLock + } + scheduledGenerators.poll()?.let { + runningGenerator = it + doSendCommandToExecute(it.shellCommand()) + generatorCommandSent.complete(Unit) + generatorCommandSent = CompletableDeferred() + } + } + } + } + + private fun Queue.drainToList(): List = ArrayList(size).also { + while (isNotEmpty()) { + it.add(poll()) + } + } + + private fun doSendCommandToExecute(shellCommand: String) { + // Simulate pressing Ctrl+U in the terminal to clear all typings in the prompt + val fullCommand = "\u0015" + shellCommand + session.terminalStarterFuture.thenAccept { + if (it != null) { + TerminalUtil.sendCommandToExecute(fullCommand, it) + } + } + } + + private class Generator(val name: String, val parameters: List) { + val requestId: Int = NEXT_REQUEST_ID.incrementAndGet() + val deferred: CompletableDeferred = CompletableDeferred() + + fun shellCommand(): String = "$name $requestId ${ParametersListUtil.join(parameters)}" + override fun toString(): String = "Generator($name, parameters=$parameters, requestId=$requestId)" + } + + private class Lock { + private val lock: Any = Any() + + fun withLock(block: (WithoutLockRegistrar) -> Unit) { + val withoutLockBlocks: MutableList<() -> Unit> = ArrayList() + try { + synchronized(lock) { + block { + withoutLockBlocks.add(it) + } + } + } + finally { + withoutLockBlocks.forEach { it() } + } + } + } + + companion object { + private val NEXT_REQUEST_ID = AtomicInteger(0) + } +} + +private typealias WithoutLockRegistrar = (() -> Unit) -> Unit \ No newline at end of file diff --git a/plugins/terminal/src/org/jetbrains/plugins/terminal/exp/ShellCommandManager.kt b/plugins/terminal/src/org/jetbrains/plugins/terminal/exp/ShellCommandManager.kt index ada7cb7d9e4a..7b5fa104c24f 100644 --- a/plugins/terminal/src/org/jetbrains/plugins/terminal/exp/ShellCommandManager.kt +++ b/plugins/terminal/src/org/jetbrains/plugins/terminal/exp/ShellCommandManager.kt @@ -3,21 +3,22 @@ package org.jetbrains.plugins.terminal.exp import com.intellij.openapi.Disposable import com.intellij.openapi.diagnostic.logger -import com.jediterm.terminal.Terminal import com.jediterm.terminal.TerminalCustomCommandListener +import kotlinx.coroutines.CompletableDeferred import org.jetbrains.plugins.terminal.TerminalUtil import java.util.* import java.util.concurrent.CopyOnWriteArrayList import java.util.concurrent.TimeUnit -class ShellCommandManager(terminal: Terminal) { +class ShellCommandManager(session: TerminalSession) { private val listeners: CopyOnWriteArrayList = CopyOnWriteArrayList() @Volatile private var startedCommand: StartedCommand? = null + internal val commandExecutionManager: ShellCommandExecutionManager = ShellCommandExecutionManager(session, this) init { - terminal.addCustomCommandListener(TerminalCustomCommandListener { + session.controller.addCustomCommandListener(TerminalCustomCommandListener { try { when (it.getOrNull(0)) { "initialized" -> processInitialized(it) @@ -150,8 +151,14 @@ class ShellCommandManager(terminal: Terminal) { TerminalUtil.addItem(listeners, listener, parentDisposable) } + fun sendCommandToExecute(shellCommand: String) = commandExecutionManager.sendCommandToExecute(shellCommand) + + fun runGeneratorAsync(generatorName: String, generatorParameters: List): CompletableDeferred { + return commandExecutionManager.runGeneratorAsync(generatorName, generatorParameters) + } + companion object { - private val LOG = logger() + internal val LOG = logger() @Throws(IllegalArgumentException::class) private fun decodeHex(hexStr: String): String { diff --git a/plugins/terminal/src/org/jetbrains/plugins/terminal/exp/TerminalSession.kt b/plugins/terminal/src/org/jetbrains/plugins/terminal/exp/TerminalSession.kt index 56d7fcdd6603..c49940c97855 100644 --- a/plugins/terminal/src/org/jetbrains/plugins/terminal/exp/TerminalSession.kt +++ b/plugins/terminal/src/org/jetbrains/plugins/terminal/exp/TerminalSession.kt @@ -3,8 +3,8 @@ package org.jetbrains.plugins.terminal.exp import com.intellij.openapi.Disposable import com.intellij.openapi.actionSystem.DataKey -import com.intellij.openapi.options.advanced.AdvancedSettings import com.intellij.openapi.diagnostic.thisLogger +import com.intellij.openapi.options.advanced.AdvancedSettings import com.intellij.openapi.util.Key import com.intellij.terminal.JBTerminalSystemSettingsProviderBase import com.intellij.terminal.TerminalColorPalette @@ -41,7 +41,7 @@ class TerminalSession(settings: JBTerminalSystemSettingsProviderBase, model = TerminalModel(textBuffer) controller = JediTerminal(ModelUpdatingTerminalDisplay(model, settings), textBuffer, styleState) - commandManager = ShellCommandManager(controller) + commandManager = ShellCommandManager(this) val typeAheadTerminalModel = JediTermTypeAheadModel(controller, textBuffer, settings) typeAheadManager = TerminalTypeAheadManager(typeAheadTerminalModel) @@ -81,13 +81,7 @@ class TerminalSession(settings: JBTerminalSystemSettingsProviderBase, } fun sendCommandToExecute(shellCommand: String) { - // Simulate pressing Ctrl+U in the terminal to clear all typings in the prompt - val fullCommand = "\u0015" + shellCommand - terminalStarterFuture.thenAccept { - if (it != null) { - TerminalUtil.sendCommandToExecute(fullCommand, it) - } - } + commandManager.sendCommandToExecute(shellCommand) } fun postResize(newSize: TermSize) { diff --git a/plugins/terminal/src/org/jetbrains/plugins/terminal/exp/completion/IJShellRuntimeDataProvider.kt b/plugins/terminal/src/org/jetbrains/plugins/terminal/exp/completion/IJShellRuntimeDataProvider.kt index 648812c5065c..bdb116adfb7a 100644 --- a/plugins/terminal/src/org/jetbrains/plugins/terminal/exp/completion/IJShellRuntimeDataProvider.kt +++ b/plugins/terminal/src/org/jetbrains/plugins/terminal/exp/completion/IJShellRuntimeDataProvider.kt @@ -3,17 +3,12 @@ package org.jetbrains.plugins.terminal.exp.completion import com.intellij.openapi.diagnostic.Logger import com.intellij.openapi.diagnostic.logger -import com.intellij.openapi.util.Disposer import com.intellij.terminal.completion.ShellEnvironment import com.intellij.terminal.completion.ShellRuntimeDataProvider -import kotlinx.coroutines.suspendCancellableCoroutine import kotlinx.serialization.Serializable import kotlinx.serialization.json.Json -import org.jetbrains.plugins.terminal.exp.ShellCommandListener import org.jetbrains.plugins.terminal.exp.TerminalSession import org.jetbrains.plugins.terminal.util.ShellType -import java.util.concurrent.atomic.AtomicInteger -import kotlin.coroutines.resume class IJShellRuntimeDataProvider(private val session: TerminalSession) : ShellRuntimeDataProvider { override suspend fun getFilesFromDirectory(path: String): List { @@ -27,30 +22,14 @@ class IJShellRuntimeDataProvider(private val session: TerminalSession) : ShellRu private suspend fun executeCommand(command: DataProviderCommand): T { return if (command.isAvailable(session)) { - val requestId = CUR_ID.getAndIncrement() - val commandText = "${command.functionName} $requestId ${command.parameters.joinToString(" ")}" - val rawResult: String = executeCommandBlocking(requestId, commandText) + val rawResult: String = executeCommandBlocking(command) command.parseResult(rawResult) } else command.defaultResult } - private suspend fun executeCommandBlocking(reqId: Int, command: String): String { - return suspendCancellableCoroutine { continuation -> - val disposable = Disposer.newDisposable() - continuation.invokeOnCancellation { - Disposer.dispose(disposable) - } - session.addCommandListener(object : ShellCommandListener { - override fun generatorFinished(requestId: Int, result: String) { - if (requestId == reqId) { - Disposer.dispose(disposable) - continuation.resume(result) - } - } - }, disposable) - session.sendCommandToExecute(command) - } + private suspend fun executeCommandBlocking(command: DataProviderCommand<*>): String { + return session.commandManager.runGeneratorAsync(command.functionName, command.parameters).await() } private interface DataProviderCommand { @@ -113,7 +92,6 @@ class IJShellRuntimeDataProvider(private val session: TerminalSession) : ShellRu } companion object { - private val CUR_ID = AtomicInteger(0) private val LOG: Logger = logger() private fun TerminalSession.isBashOrZsh(): Boolean { diff --git a/plugins/terminal/tests/org/jetbrains/plugins/terminal/block/BlockTerminalTest.kt b/plugins/terminal/tests/org/jetbrains/plugins/terminal/block/BlockTerminalTest.kt index eef870eaaf91..5529000aee4c 100644 --- a/plugins/terminal/tests/org/jetbrains/plugins/terminal/block/BlockTerminalTest.kt +++ b/plugins/terminal/tests/org/jetbrains/plugins/terminal/block/BlockTerminalTest.kt @@ -4,15 +4,21 @@ package org.jetbrains.plugins.terminal.block import com.intellij.openapi.options.advanced.AdvancedSettings import com.intellij.openapi.util.Disposer import com.intellij.openapi.util.text.StringUtil +import com.intellij.terminal.completion.ShellEnvironment import com.intellij.testFramework.DisposableRule import com.intellij.testFramework.PlatformTestUtil import com.intellij.testFramework.ProjectRule import com.intellij.testFramework.RuleChain +import com.intellij.util.TimeoutUtil import com.jediterm.core.util.TermSize +import kotlinx.coroutines.* import org.jetbrains.plugins.terminal.block.testApps.SimpleTextRepeater import org.jetbrains.plugins.terminal.exp.* +import org.jetbrains.plugins.terminal.exp.completion.IJShellRuntimeDataProvider import org.jetbrains.plugins.terminal.exp.util.TerminalSessionTestUtil +import org.jetbrains.plugins.terminal.util.ShellType import org.junit.Assert +import org.junit.Assume import org.junit.Rule import org.junit.Test import org.junit.runner.RunWith @@ -22,6 +28,8 @@ import java.util.* import java.util.concurrent.CompletableFuture import java.util.concurrent.TimeUnit import java.util.concurrent.atomic.AtomicReference +import kotlin.time.Duration.Companion.milliseconds +import kotlin.time.Duration.Companion.seconds @RunWith(Parameterized::class) class BlockTerminalTest(private val shellPath: String) { @@ -75,6 +83,30 @@ class BlockTerminalTest(private val shellPath: String) { }.attempts(1).assertTiming() } + @Test + fun `concurrent command and generator execution`() { + val session = startBlockTerminalSession() + // ShellType.FISH doesn't support generators yet + Assume.assumeTrue(setOf(ShellType.ZSH, ShellType.BASH).contains(session.shellIntegration?.shellType)) + runBlocking { + for (stepId in 1..100) { + val startNano = System.nanoTime() + val outputFuture: CompletableFuture = getCommandResultFuture(session) + val envListDeferred: Deferred = this.async(Dispatchers.Default) { + val provider = IJShellRuntimeDataProvider(session) + provider.getShellEnvironment() + } + session.commandManager.commandExecutionManager.awaitGeneratorCommandSent() + delay((1..50).random().milliseconds) // wait a little to start generator + session.sendCommandToExecute("echo foo") + val env: ShellEnvironment? = withTimeout(20.seconds) { envListDeferred.await() } + Assert.assertTrue(env != null && env.envs.isNotEmpty()) + assertCommandResult(0, "foo\n", outputFuture) + println("#$stepId Done in " + TimeoutUtil.getDurationMillis(startNano) + "ms") + } + } + } + private fun setTerminalBufferMaxLines(maxBufferLines: Int) { val terminalBufferMaxLinesCount = "terminal.buffer.max.lines.count" val prevValue = AdvancedSettings.getInt(terminalBufferMaxLinesCount) @@ -84,7 +116,7 @@ class BlockTerminalTest(private val shellPath: String) { } } - private fun startBlockTerminalSession(termSize: TermSize) = + private fun startBlockTerminalSession(termSize: TermSize = TermSize(80, 24)) = TerminalSessionTestUtil.startBlockTerminalSession(projectRule.project, shellPath, disposableRule.disposable, termSize) private fun TerminalSession.sendCommandToExecuteWithoutAddingToHistory(shellCommand: String) {