mirror of
https://gitflic.ru/project/openide/openide.git
synced 2026-09-27 10:03:11 +07:00
[terminal] IJPL-182482 Drop BackendOutputActivity
Backend output latency measurement will be reimplemented in the next commit GitOrigin-RevId: 7e541c817486a795268357d01ebcdf8e66597900
This commit is contained in:
committed by
intellij-monorepo-bot
parent
b4d528a3b0
commit
1426c4858b
+1
-6
@@ -14,7 +14,6 @@ import kotlinx.coroutines.flow.flowOf
|
||||
import kotlinx.coroutines.flow.onEach
|
||||
import org.jetbrains.plugins.terminal.block.reworked.*
|
||||
import org.jetbrains.plugins.terminal.block.ui.TerminalUiUtils
|
||||
import org.jetbrains.plugins.terminal.fus.BackendOutputActivity
|
||||
import kotlin.coroutines.cancellation.CancellationException
|
||||
|
||||
/**
|
||||
@@ -25,10 +24,7 @@ import kotlin.coroutines.cancellation.CancellationException
|
||||
* So, actually it allows restoring the state of UI that requests the [getOutputFlow].
|
||||
*/
|
||||
@OptIn(ExperimentalCoroutinesApi::class)
|
||||
internal class StateAwareTerminalSession(
|
||||
private val delegate: TerminalSession,
|
||||
private val fusActivity: BackendOutputActivity,
|
||||
) : TerminalSession {
|
||||
internal class StateAwareTerminalSession(private val delegate: TerminalSession) : TerminalSession {
|
||||
private val sessionModel: TerminalSessionModel = TerminalSessionModelImpl()
|
||||
private val outputModel: TerminalOutputModel
|
||||
private val alternateBufferModel: TerminalOutputModel
|
||||
@@ -87,7 +83,6 @@ internal class StateAwareTerminalSession(
|
||||
val styles = event.styles.map { it.toStyleRange() }
|
||||
val model = getCurrentOutputModel()
|
||||
model.updateContent(event.startLineLogicalIndex, event.text, styles)
|
||||
fusActivity.eventCollected(event)
|
||||
}
|
||||
is TerminalCursorPositionChangedEvent -> {
|
||||
val model = getCurrentOutputModel()
|
||||
|
||||
+4
-21
@@ -1,10 +1,11 @@
|
||||
package com.intellij.terminal.backend
|
||||
|
||||
import com.jediterm.core.typeahead.TerminalTypeAheadManager
|
||||
import com.jediterm.terminal.*
|
||||
import com.jediterm.terminal.emulator.JediEmulator
|
||||
import com.jediterm.terminal.TerminalDataStream
|
||||
import com.jediterm.terminal.TerminalExecutorServiceManager
|
||||
import com.jediterm.terminal.TerminalStarter
|
||||
import com.jediterm.terminal.TtyConnector
|
||||
import com.jediterm.terminal.model.JediTerminal
|
||||
import org.jetbrains.plugins.terminal.fus.BackendOutputActivity
|
||||
|
||||
internal class StopAwareTerminalStarter(
|
||||
terminal: JediTerminal,
|
||||
@@ -12,31 +13,13 @@ internal class StopAwareTerminalStarter(
|
||||
dataStream: TerminalDataStream,
|
||||
typeAheadManager: TerminalTypeAheadManager,
|
||||
executorServiceManager: TerminalExecutorServiceManager,
|
||||
private val fusActivity: BackendOutputActivity,
|
||||
) : TerminalStarter(terminal, ttyConnector, dataStream, typeAheadManager, executorServiceManager) {
|
||||
@Volatile
|
||||
var isStopped: Boolean = false
|
||||
private set
|
||||
|
||||
override fun createEmulator(dataStream: TerminalDataStream, terminal: Terminal): JediEmulator {
|
||||
return FusAwareEmulator(dataStream, terminal)
|
||||
}
|
||||
|
||||
override fun requestEmulatorStop() {
|
||||
super.requestEmulatorStop()
|
||||
isStopped = true
|
||||
}
|
||||
|
||||
// must be an inner class, because this thing is created in a superclass constructor,
|
||||
// so fusActivity is still null at that time, and we can't pass its actual value
|
||||
private inner class FusAwareEmulator(
|
||||
dataStream: TerminalDataStream,
|
||||
terminal: Terminal,
|
||||
) : JediEmulator(dataStream, terminal) {
|
||||
override fun next() {
|
||||
fusActivity.charProcessingStarted()
|
||||
super.next()
|
||||
fusActivity.charProcessingFinished()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+7
-8
@@ -11,14 +11,12 @@ import org.jetbrains.plugins.terminal.block.session.StyledCommandOutput
|
||||
import org.jetbrains.plugins.terminal.block.session.collectLines
|
||||
import org.jetbrains.plugins.terminal.block.session.scraper.SimpleStringCollector
|
||||
import org.jetbrains.plugins.terminal.block.session.scraper.StylesCollectingTerminalLinesCollector
|
||||
import org.jetbrains.plugins.terminal.fus.BackendOutputActivity
|
||||
import java.util.concurrent.CopyOnWriteArrayList
|
||||
import kotlin.math.min
|
||||
|
||||
internal class TerminalContentChangesTracker(
|
||||
private val textBuffer: TerminalTextBuffer,
|
||||
private val discardedHistoryTracker: TerminalDiscardedHistoryTracker,
|
||||
private val fusActivity: BackendOutputActivity,
|
||||
) {
|
||||
private var lastChangedVisualLine: Int = 0
|
||||
private var anyLineChanged: Boolean = false
|
||||
@@ -31,7 +29,6 @@ internal class TerminalContentChangesTracker(
|
||||
val line = textBuffer.effectiveHistoryLinesCount + fromIndex
|
||||
lastChangedVisualLine = min(lastChangedVisualLine, line)
|
||||
anyLineChanged = true
|
||||
fusActivity.processedCharsReachedTextBuffer()
|
||||
}
|
||||
|
||||
override fun linesDiscardedFromHistory(lines: List<TerminalLine>) {
|
||||
@@ -101,11 +98,13 @@ internal class TerminalContentChangesTracker(
|
||||
lastChangedVisualLine = textBuffer.effectiveHistoryLinesCount + textBuffer.screenLinesCount
|
||||
anyLineChanged = false
|
||||
|
||||
val styles = output.styleRanges.map { it.toDto() }
|
||||
val charRange = fusActivity.textBufferCharacterIndices()
|
||||
val contentUpdatedEvent = TerminalContentUpdatedEvent(output.text, styles, logicalLineIndex, charRange.first, charRange.last)
|
||||
fusActivity.textBufferCollected(contentUpdatedEvent)
|
||||
return contentUpdatedEvent
|
||||
return TerminalContentUpdatedEvent(
|
||||
text = output.text,
|
||||
styles = output.styleRanges.map { it.toDto() },
|
||||
startLineLogicalIndex = logicalLineIndex,
|
||||
firstCharIndex = -1,
|
||||
lastCharIndex = -1,
|
||||
)
|
||||
}
|
||||
|
||||
private fun scrapeOutput(startLine: Int, additionalLines: List<TerminalLine>): StyledCommandOutput {
|
||||
|
||||
@@ -12,7 +12,7 @@ import kotlinx.coroutines.channels.BufferOverflow
|
||||
import kotlinx.coroutines.flow.MutableSharedFlow
|
||||
import org.jetbrains.plugins.terminal.block.reworked.TerminalShellIntegrationEventsListener
|
||||
import org.jetbrains.plugins.terminal.block.ui.withLock
|
||||
import org.jetbrains.plugins.terminal.fus.BackendOutputActivity
|
||||
|
||||
|
||||
@OptIn(ExperimentalCoroutinesApi::class)
|
||||
internal fun createTerminalOutputFlow(
|
||||
@@ -37,7 +37,7 @@ internal fun createTerminalOutputFlow(
|
||||
)
|
||||
|
||||
val discardedHistoryTracker = TerminalDiscardedHistoryTracker(textBuffer)
|
||||
val contentChangesTracker = TerminalContentChangesTracker(textBuffer, discardedHistoryTracker, fusActivity)
|
||||
val contentChangesTracker = TerminalContentChangesTracker(textBuffer, discardedHistoryTracker)
|
||||
val cursorPositionTracker = TerminalCursorPositionTracker(textBuffer, discardedHistoryTracker, terminalDisplay)
|
||||
|
||||
/**
|
||||
|
||||
@@ -6,7 +6,6 @@ import com.intellij.openapi.project.Project
|
||||
import com.intellij.platform.util.coroutines.childScope
|
||||
import com.intellij.terminal.JBTerminalSystemSettingsProviderBase
|
||||
import com.intellij.terminal.TerminalExecutorServiceManagerImpl
|
||||
import com.intellij.terminal.backend.fus.enableFus
|
||||
import com.intellij.terminal.backend.fus.installFusListener
|
||||
import com.intellij.terminal.session.TerminalSession
|
||||
import com.intellij.terminal.session.TerminalSessionTerminatedEvent
|
||||
@@ -28,7 +27,6 @@ import kotlinx.coroutines.flow.asSharedFlow
|
||||
import kotlinx.coroutines.launch
|
||||
import org.jetbrains.plugins.terminal.LocalBlockTerminalRunner
|
||||
import org.jetbrains.plugins.terminal.ShellStartupOptions
|
||||
import org.jetbrains.plugins.terminal.fus.BackendOutputActivity
|
||||
import org.jetbrains.plugins.terminal.util.STOP_EMULATOR_TIMEOUT
|
||||
import org.jetbrains.plugins.terminal.util.waitFor
|
||||
import java.util.concurrent.CancellationException
|
||||
@@ -55,13 +53,12 @@ internal fun createTerminalSession(
|
||||
initialSize: TermSize,
|
||||
settings: JBTerminalSystemSettingsProviderBase,
|
||||
coroutineScope: CoroutineScope,
|
||||
fusActivity: BackendOutputActivity,
|
||||
): TerminalSession {
|
||||
val observableTtyConnector = ttyConnector as? ObservableTtyConnector ?: ObservableTtyConnector(ttyConnector)
|
||||
installFusListener(observableTtyConnector, fusActivity, parentDisposable = coroutineScope.asDisposable())
|
||||
installFusListener(observableTtyConnector, parentDisposable = coroutineScope.asDisposable())
|
||||
|
||||
val maxHistoryLinesCount = AdvancedSettings.getInt("terminal.buffer.max.lines.count")
|
||||
val services: JediTermServices = createJediTermServices(observableTtyConnector, fusActivity, initialSize, maxHistoryLinesCount, settings)
|
||||
val services: JediTermServices = createJediTermServices(observableTtyConnector, initialSize, maxHistoryLinesCount, settings)
|
||||
|
||||
val outputScope = coroutineScope.childScope("Terminal output forwarding")
|
||||
val shellIntegrationController = TerminalShellIntegrationController(services.controller)
|
||||
@@ -119,7 +116,7 @@ private fun createJediTermServices(
|
||||
val terminalStarter = StopAwareTerminalStarter(
|
||||
controller,
|
||||
connector,
|
||||
enableFus(TtyBasedArrayDataStream(connector), fusActivity),
|
||||
TtyBasedArrayDataStream(connector),
|
||||
typeAheadManager,
|
||||
executorService,
|
||||
fusActivity
|
||||
|
||||
@@ -10,6 +10,7 @@ import com.intellij.platform.kernel.backend.findValueEntity
|
||||
import com.intellij.platform.kernel.backend.newValueEntity
|
||||
import com.intellij.platform.util.coroutines.childScope
|
||||
import com.intellij.terminal.session.TerminalCloseEvent
|
||||
import com.intellij.terminal.session.TerminalSession
|
||||
import com.intellij.util.AwaitCancellationAndInvoke
|
||||
import com.intellij.util.awaitCancellationAndInvoke
|
||||
import com.jediterm.core.util.TermSize
|
||||
@@ -21,7 +22,6 @@ import org.jetbrains.plugins.terminal.ShellStartupOptions
|
||||
import org.jetbrains.plugins.terminal.block.reworked.session.TerminalSessionTab
|
||||
import org.jetbrains.plugins.terminal.block.reworked.session.rpc.TerminalPortForwardingId
|
||||
import org.jetbrains.plugins.terminal.block.reworked.session.rpc.TerminalSessionId
|
||||
import org.jetbrains.plugins.terminal.fus.BackendLatencyService
|
||||
import java.util.concurrent.atomic.AtomicInteger
|
||||
|
||||
@OptIn(AwaitCancellationAndInvoke::class)
|
||||
@@ -132,12 +132,10 @@ internal class TerminalTabsManager(private val project: Project, private val cor
|
||||
|
||||
val (ttyConnector, configuredOptions) = startTerminalProcess(project, optionsWithSize)
|
||||
val observableTtyConnector = ObservableTtyConnector(ttyConnector)
|
||||
val fusActivity = BackendLatencyService.getInstance().startBackendOutputActivity()
|
||||
val session = createTerminalSession(project, observableTtyConnector, termSize, JBTerminalSystemSettingsProvider(), scope, fusActivity)
|
||||
val stateAwareSession = StateAwareTerminalSession(session, fusActivity)
|
||||
val session = createTerminalSession(project, observableTtyConnector, termSize, JBTerminalSystemSettingsProvider(), scope)
|
||||
val stateAwareSession = StateAwareTerminalSession(session)
|
||||
|
||||
val sessionEntity = newValueEntity(stateAwareSession)
|
||||
fusActivity.sessionId = sessionEntity.id
|
||||
|
||||
scope.awaitCancellationAndInvoke {
|
||||
sessionEntity.delete()
|
||||
|
||||
+1
-297
@@ -1,22 +1,13 @@
|
||||
package com.intellij.terminal.backend.fus
|
||||
|
||||
import com.intellij.openapi.Disposable
|
||||
import com.intellij.openapi.diagnostic.logger
|
||||
import com.intellij.platform.rpc.UID
|
||||
import com.intellij.terminal.backend.ObservableTtyConnector
|
||||
import com.intellij.terminal.backend.TtyConnectorListener
|
||||
import com.intellij.terminal.session.TerminalContentUpdatedEvent
|
||||
import com.intellij.terminal.session.TerminalWriteBytesEvent
|
||||
import com.jediterm.terminal.TerminalDataStream
|
||||
import fleet.multiplatform.shims.ConcurrentHashMap
|
||||
import org.jetbrains.plugins.terminal.fus.BackendLatencyService
|
||||
import org.jetbrains.plugins.terminal.fus.BackendOutputActivity
|
||||
import org.jetbrains.plugins.terminal.fus.BackendTypingActivity
|
||||
import org.jetbrains.plugins.terminal.fus.ReworkedTerminalUsageCollector
|
||||
import java.util.concurrent.LinkedBlockingQueue
|
||||
import java.util.concurrent.atomic.AtomicReference
|
||||
import kotlin.time.Duration
|
||||
import kotlin.time.TimeMark
|
||||
import kotlin.time.TimeSource
|
||||
|
||||
internal class BackendLatencyServiceImpl : BackendLatencyService {
|
||||
@@ -30,22 +21,13 @@ internal class BackendLatencyServiceImpl : BackendLatencyService {
|
||||
override fun getBackendTypingActivityOrNull(bytes: ByteArray): BackendTypingActivity? {
|
||||
return backendTypingActivityByByteArray[bytes]
|
||||
}
|
||||
|
||||
override fun startBackendOutputActivity(): BackendOutputActivity {
|
||||
return BackendOutputActivityImpl()
|
||||
}
|
||||
}
|
||||
|
||||
internal fun installFusListener(
|
||||
ttyConnector: ObservableTtyConnector,
|
||||
fusActivity: BackendOutputActivity,
|
||||
parentDisposable: Disposable,
|
||||
) {
|
||||
ttyConnector.addListener(parentDisposable, object : TtyConnectorListener {
|
||||
override fun charsRead(buf: CharArray, offset: Int, length: Int) {
|
||||
fusActivity.charsRead(length)
|
||||
}
|
||||
|
||||
override fun bytesWritten(bytes: ByteArray) {
|
||||
val typingActivity = BackendLatencyService.getInstance().getBackendTypingActivityOrNull(bytes) ?: return
|
||||
try {
|
||||
@@ -58,9 +40,6 @@ internal fun installFusListener(
|
||||
})
|
||||
}
|
||||
|
||||
internal fun enableFus(stream: TerminalDataStream, fusActivity: BackendOutputActivity): TerminalDataStream =
|
||||
FusAwareTtyBasedDataStream(stream, fusActivity)
|
||||
|
||||
private val backendTypingActivityByByteArray = ConcurrentHashMap<ByteArray, BackendTypingActivityImpl>()
|
||||
|
||||
private class BackendTypingActivityImpl(override val id: Int, private val bytes: ByteArray) : BackendTypingActivity {
|
||||
@@ -77,279 +56,4 @@ private class BackendTypingActivityImpl(override val id: Int, private val bytes:
|
||||
override fun finishBytesProcessing() {
|
||||
backendTypingActivityByByteArray.remove(bytes)
|
||||
}
|
||||
}
|
||||
|
||||
private class BackendOutputActivityImpl : BackendOutputActivity {
|
||||
private val sessionIdReference = AtomicReference<UID>()
|
||||
|
||||
override var sessionId: UID?
|
||||
get() = sessionIdReference.get()
|
||||
set(value) { sessionIdReference.set(value) }
|
||||
|
||||
// split into subclasses to simplify reasoning about threads and locks
|
||||
|
||||
private val readingState = TerminalThreadReadingState()
|
||||
private val processingState = TerminalThreadProcessingState()
|
||||
private val textBufferState = TerminalThreadStateUnderTextBufferLock()
|
||||
private val eventFlowState = EventFlowState()
|
||||
|
||||
private data class ReadRange(val range: LongRange, val time: TimeMark)
|
||||
private data class LatencyPair(val min: Latency?, val max: Latency?)
|
||||
private data class Latency(val index: Long, val duration: Duration)
|
||||
|
||||
/** The part of the state that is affected by the terminal emulator thread when reading and buffering characters from the TTY. */
|
||||
private class TerminalThreadReadingState {
|
||||
/** The total number of characters read from the TTY stream and buffered. **/
|
||||
private var totalCharsRead = 0L
|
||||
/** The queue of character index ranges read from the TTY and timestamped. **/
|
||||
val readRanges = LinkedBlockingQueue<ReadRange>()
|
||||
|
||||
/** Invoked every time a new buffer is read from the TTY. */
|
||||
fun charsRead(count: Int) {
|
||||
val from = totalCharsRead
|
||||
totalCharsRead += count.toLong()
|
||||
val to = totalCharsRead
|
||||
readRanges.add(ReadRange(from until to, TimeSource.Monotonic.markNow()))
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* The part of the state that is affected by reading characters from the buffer.
|
||||
*
|
||||
* Usually accessed outside the text buffer lock, always in the terminal emulator thread.
|
||||
*/
|
||||
private class TerminalThreadProcessingState {
|
||||
/** The total number of characters read from the buffer and processed by the emulator. */
|
||||
var totalCharsProcessed = 0L
|
||||
private set
|
||||
/** The index of the first character processed during this iteration. `null` when we're not inside an iteration. */
|
||||
var thisProcessingIterationStart: Long? = null
|
||||
private set
|
||||
|
||||
/** Invoked at the start of each iteration. */
|
||||
fun charProcessingStarted() {
|
||||
thisProcessingIterationStart = totalCharsProcessed
|
||||
}
|
||||
|
||||
/** Invoked every time some characters are processed or pushed back into the buffer. In the latter case the argument is negative. */
|
||||
fun charsProcessed(count: Int) {
|
||||
totalCharsProcessed += count.toLong()
|
||||
}
|
||||
|
||||
/** Invoked at the end of each iteration. */
|
||||
fun charProcessingFinished() {
|
||||
thisProcessingIterationStart = null
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* The part of the state that is only updated or accessed under the text buffer lock.
|
||||
*
|
||||
* Not necessarily accessed from the terminal emulator thread.
|
||||
*/
|
||||
private class TerminalThreadStateUnderTextBufferLock {
|
||||
/**
|
||||
* The index of the first character that will be included in the next buffer collection event.
|
||||
* `null` if there have been no changes in the buffer since the last collection.
|
||||
*/
|
||||
private var nextTextBufferCollectionStart: Long? = null
|
||||
/**
|
||||
* The index of the last character that will be included in the next buffer collection event.
|
||||
* `null` if there have been no changes in the buffer since the last collection.
|
||||
*/
|
||||
private var nextTextBufferCollectionEnd: Long? = null
|
||||
|
||||
/**
|
||||
* Invoked every time when a processed character affects the text buffer.
|
||||
*
|
||||
* Invoked on the terminal emulator thread.
|
||||
*
|
||||
* @param processingIterationStart the index of the first character processed during this processing iteration
|
||||
* @param totalCharsProcessed the total number of processed characters, the same as the index of the last character processed so far
|
||||
*/
|
||||
fun processedCharsReachedTextBuffer(processingIterationStart: Long, totalCharsProcessed: Long) {
|
||||
// This is a bit complicated. The exact sequence is this:
|
||||
// 1. A processing iteration (com.intellij.terminal.backend.StopAwareTerminalStarter.FusAwareEmulator.next) starts.
|
||||
// 2. A character is processed.
|
||||
// 3. The text buffer may or may not be affected. If it's affected, this function is called.
|
||||
// 4. Steps 2-3 continue to repeat until the end of the iteration.
|
||||
// 5. The iteration ends.
|
||||
// Characters that don't affect the buffer are usually control characters. For example, cursor movement.
|
||||
// At any given moment the text buffer may be collected asynchronously from another thread.
|
||||
// But this collection happens under the same lock this function is invoked, so it's not THAT asynchronous.
|
||||
// The tricky part is to determine the range of character indices that match the buffer collection event.
|
||||
// We always know the number of chars already processed, but we don't know which characters actually affected the buffer.
|
||||
// We know for sure that when this callback is invoked, all characters processed so far are included in the text buffer.
|
||||
// But for the next callback, if we assume that the next range starts where the previous one ended,
|
||||
// we may end up with falsely large latencies because no-change characters from the previous iteration will be included as well.
|
||||
// To avoid this situation, we ignore the previous range end and start the range from the first character of THIS iteration.
|
||||
// But we must also account for the case when several processing iterations happen before the buffer is collected.
|
||||
// That's why we only set the range start if it wasn't set yet.
|
||||
if (nextTextBufferCollectionStart == null) {
|
||||
nextTextBufferCollectionStart = processingIterationStart
|
||||
}
|
||||
nextTextBufferCollectionEnd = totalCharsProcessed
|
||||
}
|
||||
|
||||
/**
|
||||
* Invoked every time the text buffer is collected ("scrapped").
|
||||
*
|
||||
* Invoked _not_ from the terminal emulator thread but from the collecting coroutine.
|
||||
*
|
||||
* @return the range of the characters that have made their way into the buffer since the last collection
|
||||
*/
|
||||
fun textBufferCollected(): LongRange? {
|
||||
val from = nextTextBufferCollectionStart
|
||||
val to = nextTextBufferCollectionEnd
|
||||
nextTextBufferCollectionStart = null
|
||||
nextTextBufferCollectionEnd = null
|
||||
if (from == null || to == null) {
|
||||
LOG.error("textBufferCollected, but from==$from and to==$to, both should be non-null at this point")
|
||||
return null
|
||||
}
|
||||
return from until to
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* The part of the state that is affected by collecting the text buffer and further event processing.
|
||||
*
|
||||
* Accessed from different threads, so must be thread-safe.
|
||||
*/
|
||||
private class EventFlowState {
|
||||
private val collectedRanges = ConcurrentHashMap<IdentityWrapper<TerminalContentUpdatedEvent>, LongRange>()
|
||||
|
||||
/**
|
||||
* Invoked every time the text buffer is collected ("scrapped").
|
||||
*/
|
||||
fun textBufferCollected(event: TerminalContentUpdatedEvent) {
|
||||
collectedRanges[event.toIdentity()] = event.firstCharIndex..event.lastCharIndex
|
||||
}
|
||||
|
||||
/**
|
||||
* Invoked every time the event produced by scrapping the text buffer is collected from the output flow.
|
||||
*
|
||||
* @return a pair of min/max latencies corresponding to the event char range, all non-`null` unless there's a bug somewhere
|
||||
*/
|
||||
fun eventCollected(event: TerminalContentUpdatedEvent, readRanges: LinkedBlockingQueue<ReadRange>): LatencyPair? {
|
||||
val range = collectedRanges.remove(event.toIdentity()) ?: return null
|
||||
var firstCharTime: TimeMark? = null
|
||||
var lastCharTime: TimeMark? = null
|
||||
while (true) {
|
||||
val nextRange = readRanges.peek() ?: break
|
||||
if (nextRange.range.first > range.last) break // reached the part not collected yet
|
||||
if (nextRange.range.last <= range.last) { // the entire range has been collected
|
||||
readRanges.remove()
|
||||
}
|
||||
if (range.first in nextRange.range) {
|
||||
firstCharTime = nextRange.time
|
||||
}
|
||||
if (range.last in nextRange.range) {
|
||||
lastCharTime = nextRange.time
|
||||
break
|
||||
}
|
||||
}
|
||||
// The first char will have the maximum latency, as it was sitting in the buffer the longest.
|
||||
// Compute the minimum latency first, as otherwise when they're essentially equal,
|
||||
// we can end up in a situation when the maximum is less than the minimum by some microseconds.
|
||||
val minLatency = if (lastCharTime != null) {
|
||||
Latency(range.last, lastCharTime.elapsedNow())
|
||||
}
|
||||
else {
|
||||
LOG.warn("The last char ${range.last} was lost somewhere, it's a bug")
|
||||
null
|
||||
}
|
||||
val maxLatency = if (firstCharTime != null) {
|
||||
Latency(range.first, firstCharTime.elapsedNow())
|
||||
}
|
||||
else {
|
||||
LOG.warn("The first char ${range.first} was lost somewhere, it's a bug")
|
||||
null
|
||||
}
|
||||
return LatencyPair(minLatency, maxLatency)
|
||||
}
|
||||
}
|
||||
|
||||
override fun charsRead(count: Int) = readingState.charsRead(count)
|
||||
|
||||
override fun charProcessingStarted() = processingState.charProcessingStarted()
|
||||
|
||||
override fun charsProcessed(count: Int) = processingState.charsProcessed(count)
|
||||
|
||||
override fun processedCharsReachedTextBuffer() {
|
||||
// cross-state safe interaction: this function is called in the same thread that updates processingState,
|
||||
// and under the same lock that is always used to access textBufferState
|
||||
val processingIterationStart = processingState.thisProcessingIterationStart
|
||||
if (processingIterationStart == null) {
|
||||
LOG.error("processedCharsReachedTextBuffer should not be called outside of a processing iteration")
|
||||
return
|
||||
}
|
||||
textBufferState.processedCharsReachedTextBuffer(processingIterationStart, processingState.totalCharsProcessed)
|
||||
}
|
||||
|
||||
override fun charProcessingFinished() = processingState.charProcessingFinished()
|
||||
|
||||
// cross-state safe interaction: these two functions are called under the same text buffer lock textBufferState is updated under
|
||||
|
||||
override fun textBufferCharacterIndices(): LongRange {
|
||||
return textBufferState.textBufferCollected() ?: LongRange.EMPTY
|
||||
}
|
||||
|
||||
override fun textBufferCollected(event: TerminalContentUpdatedEvent) {
|
||||
eventFlowState.textBufferCollected(event)
|
||||
}
|
||||
|
||||
override fun eventCollected(event: TerminalContentUpdatedEvent) {
|
||||
// cross-state safe interaction: using the shared queue to transfer read ranges
|
||||
val latency = eventFlowState.eventCollected(event, readingState.readRanges) ?: return
|
||||
// If the sessionId is not known yet, we still collect statistics to ensure a consistent state but skip reporting.
|
||||
// This can only happen very early during the session startup.
|
||||
val sessionId = this.sessionId ?: return
|
||||
if (latency.min != null) {
|
||||
ReworkedTerminalUsageCollector.logBackendMinOutputLatency(sessionId, latency.min.index, latency.min.duration)
|
||||
}
|
||||
if (latency.max != null) {
|
||||
ReworkedTerminalUsageCollector.logBackendMaxOutputLatency(sessionId, latency.max.index, latency.max.duration)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// used to track individual instances of data classes
|
||||
private class IdentityWrapper<T : Any>(private val instance: T) {
|
||||
override fun equals(other: Any?): Boolean = instance === (other as? IdentityWrapper<T>)?.instance
|
||||
override fun hashCode(): Int = System.identityHashCode(instance)
|
||||
}
|
||||
|
||||
private class FusAwareTtyBasedDataStream(
|
||||
private val original: TerminalDataStream,
|
||||
private val fusActivity: BackendOutputActivity,
|
||||
) : TerminalDataStream {
|
||||
override fun getChar(): Char {
|
||||
@Suppress("UsePropertyAccessSyntax")
|
||||
val result = original.getChar()
|
||||
fusActivity.charsProcessed(1)
|
||||
return result
|
||||
}
|
||||
|
||||
override fun pushChar(c: Char) {
|
||||
fusActivity.charsProcessed(-1)
|
||||
original.pushChar(c)
|
||||
}
|
||||
|
||||
override fun readNonControlCharacters(maxChars: Int): String? {
|
||||
val result = original.readNonControlCharacters(maxChars)
|
||||
fusActivity.charsProcessed(result.length)
|
||||
return result
|
||||
}
|
||||
|
||||
override fun pushBackBuffer(bytes: CharArray, length: Int) {
|
||||
fusActivity.charsProcessed(-length)
|
||||
original.pushBackBuffer(bytes, length)
|
||||
}
|
||||
|
||||
override fun isEmpty(): Boolean = original.isEmpty
|
||||
}
|
||||
|
||||
private fun <T : Any> T.toIdentity(): IdentityWrapper<T> = IdentityWrapper(this)
|
||||
|
||||
private val LOG = logger<BackendLatencyService>()
|
||||
}
|
||||
+1
-2
@@ -1,7 +1,6 @@
|
||||
// Copyright 2000-2024 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license.
|
||||
package com.intellij.terminal.backend
|
||||
|
||||
import com.intellij.terminal.backend.util.BackendOutputTestFusActivity
|
||||
import com.intellij.terminal.backend.util.scrollDown
|
||||
import com.intellij.terminal.backend.util.write
|
||||
import com.intellij.terminal.session.TerminalContentUpdatedEvent
|
||||
@@ -177,6 +176,6 @@ internal class TerminalContentChangesTrackerTest {
|
||||
|
||||
private fun createChangesTracker(textBuffer: TerminalTextBuffer): TerminalContentChangesTracker {
|
||||
val discardedHistoryTracker = TerminalDiscardedHistoryTracker(textBuffer)
|
||||
return TerminalContentChangesTracker(textBuffer, discardedHistoryTracker, BackendOutputTestFusActivity)
|
||||
return TerminalContentChangesTracker(textBuffer, discardedHistoryTracker)
|
||||
}
|
||||
}
|
||||
+1
-2
@@ -1,7 +1,6 @@
|
||||
// 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.terminal.backend
|
||||
|
||||
import com.intellij.terminal.backend.util.BackendOutputTestFusActivity
|
||||
import com.intellij.terminal.backend.util.write
|
||||
import com.intellij.terminal.session.TerminalContentUpdatedEvent
|
||||
import com.intellij.terminal.session.TerminalCursorPositionChangedEvent
|
||||
@@ -17,7 +16,7 @@ internal class TerminalCursorPositionTrackerTest {
|
||||
val textBuffer = TerminalTextBuffer(10, 5, StyleState(), 3)
|
||||
val discardedHistoryTracker = TerminalDiscardedHistoryTracker(textBuffer)
|
||||
val terminalDisplay = TerminalDisplayImpl(DefaultSettingsProvider())
|
||||
val contentChangesTracker = TerminalContentChangesTracker(textBuffer, discardedHistoryTracker, BackendOutputTestFusActivity)
|
||||
val contentChangesTracker = TerminalContentChangesTracker(textBuffer, discardedHistoryTracker)
|
||||
val cursorPositionTracker = TerminalCursorPositionTracker(textBuffer, discardedHistoryTracker, terminalDisplay)
|
||||
|
||||
// Prepare
|
||||
|
||||
-25
@@ -1,25 +0,0 @@
|
||||
package com.intellij.terminal.backend.util
|
||||
|
||||
import com.intellij.platform.rpc.UID
|
||||
import com.intellij.terminal.session.TerminalContentUpdatedEvent
|
||||
import org.jetbrains.plugins.terminal.fus.BackendOutputActivity
|
||||
|
||||
internal object BackendOutputTestFusActivity: BackendOutputActivity {
|
||||
override var sessionId: UID? = null
|
||||
|
||||
override fun charsRead(count: Int) { }
|
||||
|
||||
override fun charProcessingStarted() { }
|
||||
|
||||
override fun charsProcessed(count: Int) { }
|
||||
|
||||
override fun processedCharsReachedTextBuffer() { }
|
||||
|
||||
override fun charProcessingFinished() { }
|
||||
|
||||
override fun textBufferCharacterIndices(): LongRange = LongRange.EMPTY
|
||||
|
||||
override fun textBufferCollected(event: TerminalContentUpdatedEvent) { }
|
||||
|
||||
override fun eventCollected(event: TerminalContentUpdatedEvent) { }
|
||||
}
|
||||
+1
-3
@@ -17,7 +17,6 @@ import kotlinx.coroutines.coroutineScope
|
||||
import kotlinx.coroutines.launch
|
||||
import org.jetbrains.plugins.terminal.ShellStartupOptions
|
||||
import org.jetbrains.plugins.terminal.TerminalEngine
|
||||
import org.jetbrains.plugins.terminal.fus.BackendLatencyService
|
||||
import org.jetbrains.plugins.terminal.reworked.util.TerminalTestUtil
|
||||
import java.nio.file.Files
|
||||
import java.nio.file.Path
|
||||
@@ -36,9 +35,8 @@ internal object TerminalSessionTestUtil {
|
||||
.initialTermSize(size)
|
||||
.envVariables(mapOf(EnvironmentUtil.DISABLE_OMZ_AUTO_UPDATE to "true", "HISTFILE" to "/dev/null"))
|
||||
.build()
|
||||
val fusActivity = BackendLatencyService.getInstance().startBackendOutputActivity()
|
||||
val (ttyConnector, _) = startTerminalProcess(project, options)
|
||||
val session = createTerminalSession(project, ttyConnector, size, JBTerminalSystemSettingsProviderBase(), coroutineScope, fusActivity)
|
||||
val session = createTerminalSession(project, ttyConnector, size, JBTerminalSystemSettingsProviderBase(), coroutineScope)
|
||||
return session
|
||||
}
|
||||
|
||||
|
||||
@@ -2,8 +2,6 @@
|
||||
package org.jetbrains.plugins.terminal.fus
|
||||
|
||||
import com.intellij.openapi.components.service
|
||||
import com.intellij.platform.rpc.UID
|
||||
import com.intellij.terminal.session.TerminalContentUpdatedEvent
|
||||
import com.intellij.terminal.session.TerminalWriteBytesEvent
|
||||
import org.jetbrains.annotations.ApiStatus
|
||||
|
||||
@@ -15,7 +13,6 @@ interface BackendLatencyService {
|
||||
}
|
||||
fun tryStartBackendTypingActivity(event: TerminalWriteBytesEvent)
|
||||
fun getBackendTypingActivityOrNull(bytes: ByteArray): BackendTypingActivity?
|
||||
fun startBackendOutputActivity(): BackendOutputActivity
|
||||
}
|
||||
|
||||
@ApiStatus.Internal
|
||||
@@ -24,16 +21,3 @@ interface BackendTypingActivity {
|
||||
fun reportDuration()
|
||||
fun finishBytesProcessing()
|
||||
}
|
||||
|
||||
@ApiStatus.Internal
|
||||
interface BackendOutputActivity {
|
||||
var sessionId: UID?
|
||||
fun charsRead(count: Int)
|
||||
fun charProcessingStarted()
|
||||
fun charsProcessed(count: Int)
|
||||
fun processedCharsReachedTextBuffer()
|
||||
fun charProcessingFinished()
|
||||
fun textBufferCharacterIndices(): LongRange
|
||||
fun textBufferCollected(event: TerminalContentUpdatedEvent)
|
||||
fun eventCollected(event: TerminalContentUpdatedEvent)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user