terminal: make sure generators and user commands don't interfere with each other (IDEA-337968)

GitOrigin-RevId: dbcda6df0dd0e5dccaaf217e195aecbede173a41
This commit is contained in:
Sergey Simonchik
2023-11-29 00:04:13 +00:00
committed by intellij-monorepo-bot
parent d654500924
commit 8750ef2fe4
5 changed files with 236 additions and 39 deletions
@@ -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<Generator> = LinkedList()
private var runningGenerator: Generator? = null
private val scheduledCommands: Queue<String> = LinkedList()
private var isCommandRunning: Boolean = false
@Volatile
private var generatorCommandSent: CompletableDeferred<Unit> = 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<String>): CompletableDeferred<String> {
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 <T> Queue<T>.drainToList(): List<T> = ArrayList<T>(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<String>) {
val requestId: Int = NEXT_REQUEST_ID.incrementAndGet()
val deferred: CompletableDeferred<String> = 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
@@ -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<ShellCommandListener> = 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<String>): CompletableDeferred<String> {
return commandExecutionManager.runGeneratorAsync(generatorName, generatorParameters)
}
companion object {
private val LOG = logger<ShellCommandManager>()
internal val LOG = logger<ShellCommandManager>()
@Throws(IllegalArgumentException::class)
private fun decodeHex(hexStr: String): String {
@@ -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) {
@@ -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<String> {
@@ -27,30 +22,14 @@ class IJShellRuntimeDataProvider(private val session: TerminalSession) : ShellRu
private suspend fun <T> executeCommand(command: DataProviderCommand<T>): 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<T> {
@@ -113,7 +92,6 @@ class IJShellRuntimeDataProvider(private val session: TerminalSession) : ShellRu
}
companion object {
private val CUR_ID = AtomicInteger(0)
private val LOG: Logger = logger<IJShellRuntimeDataProvider>()
private fun TerminalSession.isBashOrZsh(): Boolean {
@@ -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<CommandResult> = getCommandResultFuture(session)
val envListDeferred: Deferred<ShellEnvironment?> = 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) {