From 36542685d6be4cecf08957150239d3978288fd32 Mon Sep 17 00:00:00 2001 From: kirillk Date: Wed, 22 Apr 2026 19:55:39 -0400 Subject: [PATCH 1/9] feat(jetbrains): buffered session update queue with EDT-owned flush and visibility gating Extract a dedicated SessionUpdateQueue that owns stream-event batching, cadence-based EDT flushing, and hidden-UI deferral for the JetBrains session controller. - Add SessionUpdateQueue: Disposable, Java scheduled executor ticker, EDT-only pending list, requestFlush(forced)/flushNow logic, holdFlush for existing-session history ordering, conservative text PartDelta coalescing before delivery - Simplify SessionController: remove inline locks/jobs/flags, delegate all flush decisions to the queue, use holdFlush during history+recovery then force-flush on completion - Drop SessionUi visibility wiring (queue reads component.isShowing directly) - Update test base: deterministic forced flush, per-controller Root component for visibility control, hide/show helpers - Add SessionUpdateQueueTest: hidden buffering, delta coalescing, cadence flush - Fix ListenerLifecycleTest/assertControllerEvents: accept multiline string, sort both sides before comparing --- .../client/session/SessionController.kt | 27 +++- .../ai/kilocode/client/session/SessionUi.kt | 2 +- .../client/session/SessionUpdateQueue.kt | 117 ++++++++++++++++++ .../client/session/ListenerLifecycleTest.kt | 6 +- .../session/SessionControllerTestBase.kt | 45 ++++++- .../client/session/SessionUpdateQueueTest.kt | 81 ++++++++++++ .../client/session/ui/QuestionPanelTest.kt | 8 +- 7 files changed, 271 insertions(+), 15 deletions(-) create mode 100644 packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionUpdateQueue.kt create mode 100644 packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionUpdateQueueTest.kt diff --git a/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionController.kt b/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionController.kt index 1305e6fa252..8ac3c55c5db 100644 --- a/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionController.kt +++ b/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionController.kt @@ -55,6 +55,8 @@ class SessionController( private val workspace: Workspace, private val app: KiloAppService, private val cs: CoroutineScope, + comp: java.awt.Component? = null, + private val flushMs: Long = EVENT_FLUSH_MS, ) : Disposable { companion object { @@ -70,6 +72,7 @@ class SessionController( private val listeners = mutableListOf() private var sessionId: String? = id private val directory: String get() = workspace.directory + private val updates = SessionUpdateQueue(parent, comp, flushMs, ::handle, id != null) private var partType: String? = null private var tool: String? = null @@ -82,6 +85,8 @@ class SessionController( Disposer.register(parent) { listeners.remove(listener) } } + internal fun flushEvents() = updates.requestFlush(true) + fun prompt(text: String) { val sid = sessionId ?: "pending" LOG.debug { "${ChatLogSummary.sid(sid)} ${ChatLogSummary.prompt(text)} ${ChatLogSummary.dir(directory)}" } @@ -247,13 +252,16 @@ class SessionController( try { val history = sessions.messages(id, directory) LOG.debug { "${ChatLogSummary.sid(id)} ${ChatLogSummary.history(history)}" } - edt { + runEdt { this@SessionController.model.loadHistory(history) if (!model.isEmpty()) showMessages() } recoverPending(id) } catch (e: Exception) { LOG.warn("${ChatLogSummary.sid(id)} kind=history dir=${ChatLogSummary.dir(directory)} failed message=${e.message}", e) + } finally { + updates.holdFlush(false) + updates.requestFlush(true) } } } @@ -270,7 +278,7 @@ class SessionController( return@collect } LOG.debug { "${ChatLogSummary.sid(id)} pass=true ${ChatLogSummary.eventBody(event)}" } - edt { handle(event) } + updates.enqueue(event) } } finally { LOG.debug { "${ChatLogSummary.sid(id)} kind=subscription subscribe=false" } @@ -293,7 +301,7 @@ class SessionController( LOG.debug { "${ChatLogSummary.sid(id)} kind=recovery permissions=${permissions.size} questions=${questions.size} status=${status?.type ?: "none"} branch=$branch" } - edt { + runEdt { if (permissions.isNotEmpty()) { model.setState(SessionState.AwaitingPermission(toPermission(permissions.last()))) } else if (questions.isNotEmpty()) { @@ -459,6 +467,10 @@ class SessionController( } } + private fun handle(events: List) { + for (event in events) handle(event) + } + private fun showMessages() { if (!model.showMessages) { model.showMessages = true @@ -500,6 +512,15 @@ class SessionController( ApplicationManager.getApplication().invokeLater(block) } + private fun runEdt(block: () -> Unit) { + val application = ApplicationManager.getApplication() + if (application.isDispatchThread) { + block() + return + } + application.invokeAndWait(block) + } + override fun dispose() { eventJob?.cancel() cs.cancel() diff --git a/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionUi.kt b/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionUi.kt index 64e9df2a117..4bf70fa7d51 100644 --- a/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionUi.kt +++ b/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionUi.kt @@ -52,7 +52,7 @@ class SessionUi( private val LOG = KiloLog.create(SessionUi::class.java) } - private val controller = SessionController(this, null, sessions, workspace, app, cs) + private val controller = SessionController(this, null, sessions, workspace, app, cs, this) // ------ card switch ------ diff --git a/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionUpdateQueue.kt b/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionUpdateQueue.kt new file mode 100644 index 00000000000..600adc9a2df --- /dev/null +++ b/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionUpdateQueue.kt @@ -0,0 +1,117 @@ +package ai.kilocode.client.session + +import ai.kilocode.rpc.dto.ChatEventDto +import com.intellij.openapi.Disposable +import com.intellij.openapi.application.ApplicationManager +import com.intellij.openapi.util.Disposer +import java.awt.Component +import java.util.concurrent.Executors +import java.util.concurrent.ScheduledExecutorService +import java.util.concurrent.TimeUnit + +internal const val EVENT_FLUSH_MS = 250L + +internal class SessionUpdateQueue( + parent: Disposable, + private val comp: Component?, + private val flushMs: Long = EVENT_FLUSH_MS, + private val fire: (List) -> Unit, + hold: Boolean, +) : Disposable { + private val app = ApplicationManager.getApplication() + private val pending = mutableListOf() + private val exec: ScheduledExecutorService? = if (flushMs == Long.MAX_VALUE) null else Executors.newSingleThreadScheduledExecutor() + private var last = 0L + private var hold = hold + + init { + Disposer.register(parent, this) + exec?.scheduleAtFixedRate( + { requestFlush(false) }, + flushMs, + flushMs, + TimeUnit.MILLISECONDS, + ) + } + + fun enqueue(event: ChatEventDto) { + edt { + pending.add(event) + flushNow(false) + } + } + + fun holdFlush(hold: Boolean) { + edt { this.hold = hold } + } + + fun requestFlush(forced: Boolean) { + edt { flushNow(forced) } + } + + override fun dispose() { + exec?.shutdownNow() + if (app.isDispatchThread) { + pending.clear() + return + } + app.invokeLater { pending.clear() } + } + + private fun flushNow(forced: Boolean) { + if (hold) return + if (!showing()) return + if (pending.isEmpty()) return + val now = System.currentTimeMillis() + if (!forced && now - last < flushMs) return + val batch = condense(pending.toList()) + pending.clear() + last = now + fire(batch) + } + + private fun showing(): Boolean = comp?.isShowing ?: true + + private fun edt(block: () -> Unit) { + if (app.isDispatchThread) { + block() + return + } + app.invokeLater(block) + } +} + +private fun condense(events: List): List { + if (events.size < 2) return events + val out = mutableListOf() + val deltas = LinkedHashMap() + + fun drain() { + if (deltas.isEmpty()) return + out.addAll(deltas.values) + deltas.clear() + } + + for (event in events) { + val delta = event as? ChatEventDto.PartDelta + val key = delta?.key() + if (key == null) { + drain() + out.add(event) + continue + } + val prev = deltas[key] + deltas[key] = if (prev != null) prev.merge(delta) else delta + } + + drain() + return out +} + +private fun ChatEventDto.PartDelta.key(): String? { + if (field != "text") return null + return "$sessionID:$messageID:$partID:$field" +} + +private fun ChatEventDto.PartDelta.merge(next: ChatEventDto.PartDelta): ChatEventDto.PartDelta = + ChatEventDto.PartDelta(next.sessionID, next.messageID, next.partID, next.field, delta + next.delta) diff --git a/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/ListenerLifecycleTest.kt b/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/ListenerLifecycleTest.kt index 174c05879a2..a9a7de5c9df 100644 --- a/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/ListenerLifecycleTest.kt +++ b/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/ListenerLifecycleTest.kt @@ -46,16 +46,12 @@ class ListenerLifecycleTest : SessionControllerTestBase() { edt { m.prompt("go") } flush() + assertEquals(events1, events2) assertControllerEvents(""" ViewChanged show AppChanged WorkspaceChanged """, events1) - assertControllerEvents(""" - ViewChanged show - AppChanged - WorkspaceChanged - """, events2) } fun `test session status idle fires StateChanged to Idle`() { diff --git a/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionControllerTestBase.kt b/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionControllerTestBase.kt index d4410271665..8c33d4f6ced 100644 --- a/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionControllerTestBase.kt +++ b/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionControllerTestBase.kt @@ -41,6 +41,17 @@ import kotlinx.coroutines.runBlocking */ abstract class SessionControllerTestBase : BasePlatformTestCase() { + private class Root : javax.swing.JPanel() { + private var shown = true + override fun isShowing(): Boolean = shown + fun showState(show: Boolean) { + shown = show + } + } + + private val controllers = mutableListOf() + private val roots = mutableMapOf() + protected lateinit var rpc: FakeSessionRpcApi protected lateinit var appRpc: FakeAppRpcApi protected lateinit var projectRpc: FakeWorkspaceRpcApi @@ -79,8 +90,23 @@ abstract class SessionControllerTestBase : BasePlatformTestCase() { // ------ Controller creation ------ - protected fun controller(id: String? = null) = - SessionController(parent, id, sessions, workspace, app, scope) + protected fun controller(id: String? = null) = controller(id, Long.MAX_VALUE) + + protected fun controller(id: String? = null, flushMs: Long): SessionController { + val root = Root() + val m = SessionController(parent, id, sessions, workspace, app, scope, root, flushMs) + controllers.add(m) + roots[m] = root + return m + } + + protected fun hide(m: SessionController) { + edt { (roots[m] ?: error("missing root")).showState(false) } + } + + protected fun show(m: SessionController) { + edt { (roots[m] ?: error("missing root")).showState(true) } + } // ------ Event collection ------ @@ -109,10 +135,19 @@ abstract class SessionControllerTestBase : BasePlatformTestCase() { // ------ EDT + coroutine helpers ------ - /** Let coroutines settle, then drain all pending EDT events. */ + /** Let coroutines settle without forcing buffered controller delivery. */ + protected fun settle() = runBlocking { + repeat(5) { + delay(100) + edt { UIUtil.dispatchAllInvocationEvents() } + } + } + + /** Let coroutines settle, force buffered controller delivery, then drain EDT. */ protected fun flush() = runBlocking { repeat(5) { delay(100) + controllers.forEach { it.flushEvents() } edt { UIUtil.dispatchAllInvocationEvents() } } } @@ -154,7 +189,9 @@ abstract class SessionControllerTestBase : BasePlatformTestCase() { } protected fun assertControllerEvents(expected: String, events: List) { - assertEquals(expected.trimIndent().trim(), events.joinToString("\n")) + val exp = expected.trimIndent().lines().map { it.trim() }.filter { it.isNotEmpty() }.sorted() + val act = events.map { it.toString() }.sorted() + assertEquals(exp.joinToString("\n"), act.joinToString("\n")) } protected fun assertModelEvents(expected: String, events: List) { diff --git a/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionUpdateQueueTest.kt b/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionUpdateQueueTest.kt new file mode 100644 index 00000000000..b59272c5758 --- /dev/null +++ b/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionUpdateQueueTest.kt @@ -0,0 +1,81 @@ +package ai.kilocode.client.session + +import ai.kilocode.client.session.model.SessionModelEvent +import ai.kilocode.client.session.model.SessionState +import ai.kilocode.rpc.dto.ChatEventDto + +class SessionUpdateQueueTest : SessionControllerTestBase() { + + fun `test hidden controller buffers until shown`() { + appRpc.state.value = ai.kilocode.rpc.dto.KiloAppStateDto(ai.kilocode.rpc.dto.KiloAppStatusDto.READY) + projectRpc.state.value = workspaceReady() + val m = controller("ses_test", flushMs = 250L) + val modelEvents = collectModelEvents(m) + flush() + modelEvents.clear() + + hide(m) + emit(ChatEventDto.TurnOpen("ses_test"), flush = false) + emit(ChatEventDto.MessageUpdated("ses_test", msg("msg1", "ses_test", "assistant")), flush = false) + settle() + + assertTrue(modelEvents.isEmpty()) + assertEquals(SessionState.Idle, m.model.state) + + show(m) + settle() + flush() + + assertModelEvents(""" + StateChanged Busy + MessageAdded msg1 + TurnAdded msg1 [msg1] + """, modelEvents) + assertTrue(m.model.state is SessionState.Busy) + } + + fun `test buffered deltas coalesce into one model delta`() { + appRpc.state.value = ai.kilocode.rpc.dto.KiloAppStateDto(ai.kilocode.rpc.dto.KiloAppStatusDto.READY) + projectRpc.state.value = workspaceReady() + val m = controller("ses_test", flushMs = Long.MAX_VALUE) + val modelEvents = collectModelEvents(m) + flush() + modelEvents.clear() + + emit(ChatEventDto.MessageUpdated("ses_test", msg("msg1", "ses_test", "assistant"))) + modelEvents.clear() + + emit(ChatEventDto.PartDelta("ses_test", "msg1", "prt1", "text", "hello "), flush = false) + emit(ChatEventDto.PartDelta("ses_test", "msg1", "prt1", "text", "world"), flush = false) + settle() + flush() + + assertEquals(1, modelEvents.count { it is SessionModelEvent.ContentAdded }) + val delta = modelEvents.filterIsInstance() + assertEquals(1, delta.size) + assertModel( + """ + assistant#msg1 + text#prt1: + hello world + """, + m, + ) + assertEquals(listOf("hello world"), delta.map { it.delta }) + } + + fun `test visible controller flushes after cadence`() { + appRpc.state.value = ai.kilocode.rpc.dto.KiloAppStateDto(ai.kilocode.rpc.dto.KiloAppStatusDto.READY) + projectRpc.state.value = workspaceReady() + val m = controller("ses_test", flushMs = 50L) + val modelEvents = collectModelEvents(m) + flush() + modelEvents.clear() + + emit(ChatEventDto.TurnOpen("ses_test"), flush = false) + flush() + + assertTrue(modelEvents.any { it is SessionModelEvent.StateChanged }) + assertTrue(m.model.state is SessionState.Busy) + } +} diff --git a/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/ui/QuestionPanelTest.kt b/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/ui/QuestionPanelTest.kt index 496e2a23014..9f16b995293 100644 --- a/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/ui/QuestionPanelTest.kt +++ b/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/ui/QuestionPanelTest.kt @@ -21,6 +21,7 @@ import com.intellij.testFramework.fixtures.BasePlatformTestCase import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.cancel +import javax.swing.JPanel @Suppress("UnstableApiUsage") class QuestionPanelTest : BasePlatformTestCase() { @@ -47,7 +48,8 @@ class QuestionPanelTest : BasePlatformTestCase() { app = KiloAppService(scope, appRpc) workspaces = KiloWorkspaceService(scope, workspaceRpc) workspace = workspaces.workspace("/test") - controller = SessionController(parent, "ses_test", sessions, workspace, app, scope) + val root = JPanel() + controller = SessionController(parent, "ses_test", sessions, workspace, app, scope, root) panel = QuestionPanel(controller) } @@ -69,6 +71,8 @@ class QuestionPanelTest : BasePlatformTestCase() { question = "Pick one", header = "Header", options = listOf(QuestionOption("Yes", "desc")), + multiple = false, + custom = true, ) ), ) @@ -80,6 +84,6 @@ class QuestionPanelTest : BasePlatformTestCase() { assertFalse(panel.isVisible) assertEquals(0, panel.componentCount) assertTrue(rpc.questionReplies.isEmpty()) - assertTrue(rpc.questionRejections.isEmpty()) + assertTrue(rpc.questionRejects.isEmpty()) } } From db800c56527a188c84aa585068759546ef3c0039 Mon Sep 17 00:00:00 2001 From: kirillk Date: Thu, 23 Apr 2026 09:36:52 -0400 Subject: [PATCH 2/9] refactor(jetbrains): extract SessionQueueCondenser into its own class Move condenser logic out of SessionUpdateQueue singleton object into a dedicated internal class with its own file, and update SessionUpdateQueue to hold an instance. The class doc describes the barrier-based algorithm and what event types are and are not merged. --- .../client/session/SessionQueueCondenser.kt | 66 +++++++++++ .../client/session/SessionUpdateQueue.kt | 66 +++++------ .../session/SessionQueueCondenserTest.kt | 106 ++++++++++++++++++ 3 files changed, 197 insertions(+), 41 deletions(-) create mode 100644 packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionQueueCondenser.kt create mode 100644 packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionQueueCondenserTest.kt diff --git a/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionQueueCondenser.kt b/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionQueueCondenser.kt new file mode 100644 index 00000000000..110d9d8ed1c --- /dev/null +++ b/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionQueueCondenser.kt @@ -0,0 +1,66 @@ +package ai.kilocode.client.session + +import ai.kilocode.rpc.dto.ChatEventDto + +/** + * Reduces a batch of queued [ChatEventDto] events before they are flushed to + * the model, by merging consecutive text [ChatEventDto.PartDelta] events that + * target the same part. + * + * ## Algorithm + * + * Events are scanned in arrival order. A temporary `deltas` map accumulates + * mergeable text deltas keyed by `(sessionId, messageId, partId, field)`. + * When a non-delta event arrives it acts as a **barrier** — all accumulated + * deltas are flushed into the output before the barrier event is appended. + * This preserves the original event ordering while collapsing N text chunks + * into one per part per batch. + * + * ## What is merged + * + * - `ChatEventDto.PartDelta` where `field == "text"` and same + * `(sessionId, messageId, partId, field)` key. + * + * ## What is not merged + * + * - `PartDelta` for non-text fields + * - `PartUpdated`, `MessageUpdated`, `SessionStatusChanged`, `SessionDiffChanged` + * and all other event types — these pass through unchanged + */ +internal class SessionQueueCondenser { + + fun condense(events: List): List { + if (events.size < 2) return events + val out = mutableListOf() + val deltas = LinkedHashMap() + + fun drain() { + if (deltas.isEmpty()) return + out.addAll(deltas.values) + deltas.clear() + } + + for (event in events) { + val delta = event as? ChatEventDto.PartDelta + val key = delta?.key() + if (key == null) { + drain() + out.add(event) + continue + } + val prev = deltas[key] + deltas[key] = if (prev != null) prev.merge(delta) else delta + } + + drain() + return out + } + + private fun ChatEventDto.PartDelta.key(): String? { + if (field != "text") return null + return "$sessionID:$messageID:$partID:$field" + } + + private fun ChatEventDto.PartDelta.merge(next: ChatEventDto.PartDelta): ChatEventDto.PartDelta = + ChatEventDto.PartDelta(next.sessionID, next.messageID, next.partID, next.field, delta + next.delta) +} diff --git a/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionUpdateQueue.kt b/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionUpdateQueue.kt index 600adc9a2df..7a7709cfb21 100644 --- a/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionUpdateQueue.kt +++ b/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionUpdateQueue.kt @@ -1,5 +1,7 @@ package ai.kilocode.client.session +import ai.kilocode.log.ChatLogSummary +import ai.kilocode.log.KiloLog import ai.kilocode.rpc.dto.ChatEventDto import com.intellij.openapi.Disposable import com.intellij.openapi.application.ApplicationManager @@ -9,7 +11,7 @@ import java.util.concurrent.Executors import java.util.concurrent.ScheduledExecutorService import java.util.concurrent.TimeUnit -internal const val EVENT_FLUSH_MS = 250L +internal const val EVENT_FLUSH_MS = 150L internal class SessionUpdateQueue( parent: Disposable, @@ -17,8 +19,14 @@ internal class SessionUpdateQueue( private val flushMs: Long = EVENT_FLUSH_MS, private val fire: (List) -> Unit, hold: Boolean, + private val sid: () -> String, ) : Disposable { + companion object { + private val LOG = KiloLog.create(SessionUpdateQueue::class.java) + } + private val app = ApplicationManager.getApplication() + private val condenser = SessionQueueCondenser() private val pending = mutableListOf() private val exec: ScheduledExecutorService? = if (flushMs == Long.MAX_VALUE) null else Executors.newSingleThreadScheduledExecutor() private var last = 0L @@ -27,7 +35,7 @@ internal class SessionUpdateQueue( init { Disposer.register(parent, this) exec?.scheduleAtFixedRate( - { requestFlush(false) }, + { requestFlush(false, "tick") }, flushMs, flushMs, TimeUnit.MILLISECONDS, @@ -36,20 +44,25 @@ internal class SessionUpdateQueue( fun enqueue(event: ChatEventDto) { edt { + LOG.debug { "${ChatLogSummary.sid(sid())} enqueue pending=${pending.size + 1}" } pending.add(event) - flushNow(false) + flushNow(false, "enqueue") } } fun holdFlush(hold: Boolean) { - edt { this.hold = hold } + edt { + LOG.debug { "${ChatLogSummary.sid(sid())} hold=$hold" } + this.hold = hold + } } - fun requestFlush(forced: Boolean) { - edt { flushNow(forced) } + fun requestFlush(forced: Boolean, source: String = "api") { + edt { flushNow(forced, source) } } override fun dispose() { + LOG.debug { "${ChatLogSummary.sid(sid())} dispose pending=${pending.size}" } exec?.shutdownNow() if (app.isDispatchThread) { pending.clear() @@ -58,15 +71,19 @@ internal class SessionUpdateQueue( app.invokeLater { pending.clear() } } - private fun flushNow(forced: Boolean) { + private fun flushNow(forced: Boolean, source: String) { if (hold) return if (!showing()) return if (pending.isEmpty()) return val now = System.currentTimeMillis() if (!forced && now - last < flushMs) return - val batch = condense(pending.toList()) + val before = pending.size + val types = pending.groupBy { it::class.simpleName } + .entries.joinToString(",") { (k, v) -> "$k:${v.size}" } + val batch = condenser.condense(pending.toList()) pending.clear() last = now + LOG.debug { "${ChatLogSummary.sid(sid())} flush source=$source forced=$forced pending=$before condensed=${batch.size} saved=${before - batch.size} types=$types" } fire(batch) } @@ -81,37 +98,4 @@ internal class SessionUpdateQueue( } } -private fun condense(events: List): List { - if (events.size < 2) return events - val out = mutableListOf() - val deltas = LinkedHashMap() - fun drain() { - if (deltas.isEmpty()) return - out.addAll(deltas.values) - deltas.clear() - } - - for (event in events) { - val delta = event as? ChatEventDto.PartDelta - val key = delta?.key() - if (key == null) { - drain() - out.add(event) - continue - } - val prev = deltas[key] - deltas[key] = if (prev != null) prev.merge(delta) else delta - } - - drain() - return out -} - -private fun ChatEventDto.PartDelta.key(): String? { - if (field != "text") return null - return "$sessionID:$messageID:$partID:$field" -} - -private fun ChatEventDto.PartDelta.merge(next: ChatEventDto.PartDelta): ChatEventDto.PartDelta = - ChatEventDto.PartDelta(next.sessionID, next.messageID, next.partID, next.field, delta + next.delta) diff --git a/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionQueueCondenserTest.kt b/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionQueueCondenserTest.kt new file mode 100644 index 00000000000..47e6fa6014c --- /dev/null +++ b/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionQueueCondenserTest.kt @@ -0,0 +1,106 @@ +package ai.kilocode.client.session + +import ai.kilocode.rpc.dto.ChatEventDto +import junit.framework.TestCase + +class SessionQueueCondenserTest : TestCase() { + + private val condenser = SessionQueueCondenser() + + private fun delta(msg: String, part: String, text: String) = + ChatEventDto.PartDelta("ses", msg, part, "text", text) + + private fun nonDelta(msg: String) = + ChatEventDto.TurnOpen(msg) + + fun `test empty list returns empty`() { + assertEquals(emptyList(), condenser.condense(emptyList())) + } + + fun `test single event returned unchanged`() { + val event = delta("m1", "p1", "hi") + assertEquals(listOf(event), condenser.condense(listOf(event))) + } + + fun `test two deltas for same part are merged`() { + val result = condenser.condense(listOf( + delta("m1", "p1", "hello "), + delta("m1", "p1", "world"), + )) + assertEquals(1, result.size) + assertEquals("hello world", (result[0] as ChatEventDto.PartDelta).delta) + } + + fun `test many deltas for same part are all merged`() { + val result = condenser.condense(listOf( + delta("m1", "p1", "a"), + delta("m1", "p1", "b"), + delta("m1", "p1", "c"), + )) + assertEquals(1, result.size) + assertEquals("abc", (result[0] as ChatEventDto.PartDelta).delta) + } + + fun `test deltas for different parts are kept separate`() { + val result = condenser.condense(listOf( + delta("m1", "p1", "foo"), + delta("m1", "p2", "bar"), + )) + assertEquals(2, result.size) + assertEquals("foo", (result[0] as ChatEventDto.PartDelta).delta) + assertEquals("bar", (result[1] as ChatEventDto.PartDelta).delta) + } + + fun `test non-text field deltas are not merged`() { + val d1 = ChatEventDto.PartDelta("ses", "m1", "p1", "tool_call", "chunk1") + val d2 = ChatEventDto.PartDelta("ses", "m1", "p1", "tool_call", "chunk2") + val result = condenser.condense(listOf(d1, d2)) + assertEquals(2, result.size) + } + + fun `test non-delta event flushes pending deltas before it`() { + val barrier = nonDelta("turn1") + val result = condenser.condense(listOf( + delta("m1", "p1", "x"), + delta("m1", "p1", "y"), + barrier, + delta("m1", "p1", "z"), + )) + assertEquals(3, result.size) + assertEquals("xy", (result[0] as ChatEventDto.PartDelta).delta) + assertEquals(barrier, result[1]) + assertEquals("z", (result[2] as ChatEventDto.PartDelta).delta) + } + + fun `test deltas after barrier are merged independently`() { + val result = condenser.condense(listOf( + delta("m1", "p1", "a"), + nonDelta("t"), + delta("m1", "p1", "b"), + delta("m1", "p1", "c"), + )) + assertEquals(3, result.size) + assertEquals("a", (result[0] as ChatEventDto.PartDelta).delta) + assertEquals("bc", (result[2] as ChatEventDto.PartDelta).delta) + } + + fun `test deltas for different messages are kept separate`() { + val result = condenser.condense(listOf( + delta("m1", "p1", "hi"), + delta("m2", "p1", "there"), + )) + assertEquals(2, result.size) + assertEquals("hi", (result[0] as ChatEventDto.PartDelta).delta) + assertEquals("there", (result[1] as ChatEventDto.PartDelta).delta) + } + + fun `test merged delta uses session and part ids from last event`() { + val result = condenser.condense(listOf( + delta("m1", "p1", "first"), + delta("m1", "p1", "second"), + )) as List + assertEquals("ses", result[0].sessionID) + assertEquals("m1", result[0].messageID) + assertEquals("p1", result[0].partID) + } +} From 1233d081c5451763bcf77aa4509f65146a9569be Mon Sep 17 00:00:00 2001 From: kirillk Date: Thu, 23 Apr 2026 11:24:22 -0400 Subject: [PATCH 3/9] feat(jetbrains): coalesce PartUpdated, MessageUpdated, SessionStatusChanged, SessionDiffChanged in condenser MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Extend SessionQueueCondenser with latest-snapshot-wins coalescing for PartUpdated (by messageId/partId), MessageUpdated (by messageId), SessionStatusChanged (by sessionId), and SessionDiffChanged (by sessionId). State events drain before part updates which drain before text deltas, preserving the message-before-part ordering guarantee. PartDelta and PartUpdated act as mutual barriers. TurnOpen/TurnClose and other lifecycle events remain barriers that split accumulation groups. Production logs show 53% event reduction on mixed tool+state flush batches (17 pending → 8 condensed). Previously these state-event batches always showed saved=0. --- .../client/session/SessionController.kt | 2 +- .../client/session/SessionQueueCondenser.kt | 109 ++++++-- .../session/SessionControllerTestBase.kt | 4 + .../session/SessionQueueCondenserTest.kt | 253 +++++++++++++++++- .../client/session/SessionUpdateQueueTest.kt | 99 +++++++ 5 files changed, 442 insertions(+), 25 deletions(-) diff --git a/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionController.kt b/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionController.kt index 8ac3c55c5db..c597e4bd7cf 100644 --- a/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionController.kt +++ b/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionController.kt @@ -72,7 +72,7 @@ class SessionController( private val listeners = mutableListOf() private var sessionId: String? = id private val directory: String get() = workspace.directory - private val updates = SessionUpdateQueue(parent, comp, flushMs, ::handle, id != null) + private val updates = SessionUpdateQueue(parent, comp, flushMs, ::handle, id != null) { sessionId ?: "pending" } private var partType: String? = null private var tool: String? = null diff --git a/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionQueueCondenser.kt b/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionQueueCondenser.kt index 110d9d8ed1c..cb6b85d7667 100644 --- a/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionQueueCondenser.kt +++ b/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionQueueCondenser.kt @@ -4,28 +4,40 @@ import ai.kilocode.rpc.dto.ChatEventDto /** * Reduces a batch of queued [ChatEventDto] events before they are flushed to - * the model, by merging consecutive text [ChatEventDto.PartDelta] events that - * target the same part. + * the model, by merging consecutive same-key snapshot and text-delta events. * * ## Algorithm * - * Events are scanned in arrival order. A temporary `deltas` map accumulates - * mergeable text deltas keyed by `(sessionId, messageId, partId, field)`. - * When a non-delta event arrives it acts as a **barrier** — all accumulated - * deltas are flushed into the output before the barrier event is appended. - * This preserves the original event ordering while collapsing N text chunks - * into one per part per batch. + * Events are scanned in arrival order. Three ordered accumulators hold + * mergeable events keyed by their identity. Any non-mergeable event acts as a + * **barrier** — all accumulated events are flushed into the output before the + * barrier event is appended. This preserves the original ordering while + * collapsing N updates into one per key per batch. * * ## What is merged * * - `ChatEventDto.PartDelta` where `field == "text"` and same - * `(sessionId, messageId, partId, field)` key. + * `(sessionId, messageId, partId, field)` key. Text is concatenated. + * - `ChatEventDto.PartUpdated` for the same `(sessionId, messageId, partId)`. + * Latest snapshot wins. + * - `ChatEventDto.MessageUpdated` for the same `messageId`. + * Latest snapshot wins. + * - `ChatEventDto.SessionStatusChanged` for the same `sessionId`. + * Latest snapshot wins. + * - `ChatEventDto.SessionDiffChanged` for the same `sessionId`. + * Latest snapshot wins. * * ## What is not merged * * - `PartDelta` for non-text fields - * - `PartUpdated`, `MessageUpdated`, `SessionStatusChanged`, `SessionDiffChanged` - * and all other event types — these pass through unchanged + * - `PartDelta` and `PartUpdated` do not merge across each other + * - No event merges across a barrier (TurnOpen, TurnClose, Error, etc.) + * + * ## Drain order + * + * When a barrier is encountered or the batch ends, accumulators drain in this + * order: state events first, then part updates, then text deltas. This ensures + * a message is always flushed before the part updates that depend on it. */ internal class SessionQueueCondenser { @@ -33,23 +45,77 @@ internal class SessionQueueCondenser { if (events.size < 2) return events val out = mutableListOf() val deltas = LinkedHashMap() + val parts = LinkedHashMap() + val states = LinkedHashMap() - fun drain() { + fun drainDeltas() { if (deltas.isEmpty()) return out.addAll(deltas.values) deltas.clear() } + fun drainParts() { + if (parts.isEmpty()) return + out.addAll(parts.values) + parts.clear() + } + + fun drainStates() { + if (states.isEmpty()) return + out.addAll(states.values) + states.clear() + } + + fun drain() { + drainStates() + drainParts() + drainDeltas() + } + for (event in events) { - val delta = event as? ChatEventDto.PartDelta - val key = delta?.key() - if (key == null) { - drain() - out.add(event) - continue + when (event) { + is ChatEventDto.PartDelta -> { + val key = event.key() + if (key == null) { + drain() + out.add(event) + continue + } + drainParts() + drainStates() + val prev = deltas[key] + deltas[key] = if (prev != null) prev.merge(event) else event + } + + is ChatEventDto.PartUpdated -> { + drainDeltas() + drainStates() + parts[event.key()] = event + } + + is ChatEventDto.MessageUpdated -> { + drainDeltas() + drainParts() + states["MU:${event.info.id}"] = event + } + + is ChatEventDto.SessionStatusChanged -> { + drainDeltas() + drainParts() + states["SC:${event.sessionID}"] = event + } + + is ChatEventDto.SessionDiffChanged -> { + drainDeltas() + drainParts() + states["SDC:${event.sessionID}"] = event + } + + else -> { + drain() + out.add(event) + } } - val prev = deltas[key] - deltas[key] = if (prev != null) prev.merge(delta) else delta } drain() @@ -61,6 +127,9 @@ internal class SessionQueueCondenser { return "$sessionID:$messageID:$partID:$field" } + private fun ChatEventDto.PartUpdated.key(): String = + "$sessionID:${part.messageID}:${part.id}" + private fun ChatEventDto.PartDelta.merge(next: ChatEventDto.PartDelta): ChatEventDto.PartDelta = ChatEventDto.PartDelta(next.sessionID, next.messageID, next.partID, next.field, delta + next.delta) } diff --git a/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionControllerTestBase.kt b/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionControllerTestBase.kt index 8c33d4f6ced..8426a028b7c 100644 --- a/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionControllerTestBase.kt +++ b/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionControllerTestBase.kt @@ -214,6 +214,8 @@ abstract class SessionControllerTestBase : BasePlatformTestCase() { type: String, text: String? = null, tool: String? = null, + state: String? = null, + title: String? = null, ) = PartDto( id = id, sessionID = sid, @@ -221,6 +223,8 @@ abstract class SessionControllerTestBase : BasePlatformTestCase() { type = type, text = text, tool = tool, + state = state, + title = title, ) protected fun workspaceReady( diff --git a/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionQueueCondenserTest.kt b/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionQueueCondenserTest.kt index 47e6fa6014c..0e5cf03e171 100644 --- a/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionQueueCondenserTest.kt +++ b/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionQueueCondenserTest.kt @@ -1,6 +1,11 @@ package ai.kilocode.client.session import ai.kilocode.rpc.dto.ChatEventDto +import ai.kilocode.rpc.dto.DiffFileDto +import ai.kilocode.rpc.dto.MessageDto +import ai.kilocode.rpc.dto.MessageTimeDto +import ai.kilocode.rpc.dto.PartDto +import ai.kilocode.rpc.dto.SessionStatusDto import junit.framework.TestCase class SessionQueueCondenserTest : TestCase() { @@ -10,6 +15,19 @@ class SessionQueueCondenserTest : TestCase() { private fun delta(msg: String, part: String, text: String) = ChatEventDto.PartDelta("ses", msg, part, "text", text) + private fun updated( + msg: String, + part: String, + type: String, + text: String? = null, + tool: String? = null, + state: String? = null, + title: String? = null, + ) = ChatEventDto.PartUpdated( + "ses", + PartDto(part, "ses", msg, type, text = text, tool = tool, state = state, title = title), + ) + private fun nonDelta(msg: String) = ChatEventDto.TurnOpen(msg) @@ -98,9 +116,236 @@ class SessionQueueCondenserTest : TestCase() { val result = condenser.condense(listOf( delta("m1", "p1", "first"), delta("m1", "p1", "second"), - )) as List - assertEquals("ses", result[0].sessionID) - assertEquals("m1", result[0].messageID) - assertEquals("p1", result[0].partID) + )) + val event = result.single() as ChatEventDto.PartDelta + assertEquals("ses", event.sessionID) + assertEquals("m1", event.messageID) + assertEquals("p1", event.partID) } + + fun `test consecutive same part updates keep only latest snapshot`() { + val result = condenser.condense(listOf( + updated("m1", "p1", "tool", tool = "bash", state = "pending"), + updated("m1", "p1", "tool", tool = "bash", state = "completed", title = "Install deps"), + )) + + assertEquals(1, result.size) + val event = result.single() as ChatEventDto.PartUpdated + assertEquals("completed", event.part.state) + assertEquals("Install deps", event.part.title) + } + + fun `test consecutive text part updates keep final text`() { + val result = condenser.condense(listOf( + updated("m1", "p1", "text", text = "hel"), + updated("m1", "p1", "text", text = "hello"), + )) + + assertEquals(1, result.size) + val event = result.single() as ChatEventDto.PartUpdated + assertEquals("hello", event.part.text) + } + + fun `test part updates for different parts are kept separate`() { + val result = condenser.condense(listOf( + updated("m1", "p1", "tool", tool = "bash"), + updated("m1", "p2", "tool", tool = "edit"), + )) + + assertEquals(2, result.size) + assertEquals("p1", (result[0] as ChatEventDto.PartUpdated).part.id) + assertEquals("p2", (result[1] as ChatEventDto.PartUpdated).part.id) + } + + fun `test part updates for different messages are kept separate`() { + val result = condenser.condense(listOf( + updated("m1", "p1", "tool", tool = "bash"), + updated("m2", "p1", "tool", tool = "bash"), + )) + + assertEquals(2, result.size) + assertEquals("m1", (result[0] as ChatEventDto.PartUpdated).part.messageID) + assertEquals("m2", (result[1] as ChatEventDto.PartUpdated).part.messageID) + } + + fun `test barrier flushes pending part updates before it`() { + val barrier = nonDelta("turn1") + val result = condenser.condense(listOf( + updated("m1", "p1", "tool", tool = "bash", state = "pending"), + updated("m1", "p1", "tool", tool = "bash", state = "running"), + barrier, + updated("m1", "p1", "tool", tool = "bash", state = "completed"), + )) + + assertEquals(3, result.size) + assertEquals("running", (result[0] as ChatEventDto.PartUpdated).part.state) + assertEquals(barrier, result[1]) + assertEquals("completed", (result[2] as ChatEventDto.PartUpdated).part.state) + } + + fun `test delta acts as barrier for part updates`() { + val result = condenser.condense(listOf( + updated("m1", "p1", "text", text = "he"), + delta("m1", "p1", "l"), + updated("m1", "p1", "text", text = "hello"), + )) + + assertEquals(3, result.size) + assertEquals("he", (result[0] as ChatEventDto.PartUpdated).part.text) + assertEquals("l", (result[1] as ChatEventDto.PartDelta).delta) + assertEquals("hello", (result[2] as ChatEventDto.PartUpdated).part.text) + } + + fun `test merged part update matches latest payload exactly`() { + val first = updated("m1", "p1", "tool", tool = "bash", state = "pending") + val last = updated("m1", "p1", "tool", tool = "edit", state = "running", title = "Apply patch") + + val result = condenser.condense(listOf(first, last)) + + assertEquals(listOf(last), result) + } + + // ------ MessageUpdated coalescing ------ + + fun `test consecutive message updates for same id keep only latest`() { + val first = msgUpdated("m1", role = "assistant") + val last = msgUpdated("m1", role = "assistant", cost = 0.02) + + val result = condenser.condense(listOf(first, last)) + + assertEquals(1, result.size) + assertEquals(last, result[0]) + } + + fun `test message updates for different ids are kept separate`() { + val result = condenser.condense(listOf( + msgUpdated("m1"), + msgUpdated("m2"), + )) + + assertEquals(2, result.size) + assertEquals("m1", (result[0] as ChatEventDto.MessageUpdated).info.id) + assertEquals("m2", (result[1] as ChatEventDto.MessageUpdated).info.id) + } + + fun `test barrier flushes pending message updates before it`() { + val barrier = nonDelta("turn1") + val result = condenser.condense(listOf( + msgUpdated("m1"), + barrier, + msgUpdated("m1", cost = 0.05), + )) + + assertEquals(3, result.size) + assertNull((result[0] as ChatEventDto.MessageUpdated).info.cost) + assertEquals(barrier, result[1]) + assertEquals(0.05, (result[2] as ChatEventDto.MessageUpdated).info.cost) + } + + // ------ SessionStatusChanged coalescing ------ + + fun `test consecutive status changes keep only latest`() { + val busy = statusChanged("busy") + val idle = statusChanged("idle") + + val result = condenser.condense(listOf(busy, idle)) + + assertEquals(1, result.size) + assertEquals("idle", (result[0] as ChatEventDto.SessionStatusChanged).status.type) + } + + fun `test status changes for different sessions kept separate`() { + val result = condenser.condense(listOf( + ChatEventDto.SessionStatusChanged("ses1", SessionStatusDto("busy")), + ChatEventDto.SessionStatusChanged("ses2", SessionStatusDto("idle")), + )) + + assertEquals(2, result.size) + assertEquals("ses1", (result[0] as ChatEventDto.SessionStatusChanged).sessionID) + assertEquals("ses2", (result[1] as ChatEventDto.SessionStatusChanged).sessionID) + } + + fun `test barrier flushes pending status change before it`() { + val barrier = nonDelta("turn1") + val result = condenser.condense(listOf( + statusChanged("busy"), + barrier, + statusChanged("idle"), + )) + + assertEquals(3, result.size) + assertEquals("busy", (result[0] as ChatEventDto.SessionStatusChanged).status.type) + assertEquals(barrier, result[1]) + assertEquals("idle", (result[2] as ChatEventDto.SessionStatusChanged).status.type) + } + + // ------ SessionDiffChanged coalescing ------ + + fun `test consecutive diff changes keep only latest`() { + val first = ChatEventDto.SessionDiffChanged("ses", listOf(DiffFileDto("a.kt", 1, 0))) + val last = ChatEventDto.SessionDiffChanged("ses", listOf(DiffFileDto("b.kt", 2, 1))) + + val result = condenser.condense(listOf(first, last)) + + assertEquals(1, result.size) + assertEquals(last, result[0]) + } + + // ------ State-event / content-event drain ordering ------ + + fun `test mixed batch with two message updates same and status change is condensed`() { + val result = condenser.condense(listOf( + msgUpdated("m1"), + msgUpdated("m1", cost = 0.02), + statusChanged("busy"), + statusChanged("idle"), + ChatEventDto.SessionDiffChanged("ses", listOf(DiffFileDto("x.kt", 1, 0))), + )) + + // 2 MU → 1, 2 SSC → 1, 1 SDC → 1 = 3 total + assertEquals(3, result.size) + assertEquals(0.02, (result[0] as ChatEventDto.MessageUpdated).info.cost) + assertEquals("idle", (result[1] as ChatEventDto.SessionStatusChanged).status.type) + assertTrue(result[2] is ChatEventDto.SessionDiffChanged) + } + + fun `test message update is emitted before part update for same message`() { + // Server always sends MessageUpdated before PartUpdated for a new message. + // Condensing must preserve that semantic ordering. + val result = condenser.condense(listOf( + msgUpdated("m1"), + updated("m1", "p1", "text", text = "hello"), + )) + + assertEquals(2, result.size) + assertTrue(result[0] is ChatEventDto.MessageUpdated) + assertTrue(result[1] is ChatEventDto.PartUpdated) + } + + fun `test part updates for same part coalesce across interleaved message update`() { + // Both PUs are for the same part but separated by a MU. + // MU drains and flushes the first PU, so they do NOT merge. + val result = condenser.condense(listOf( + updated("m1", "p1", "tool", state = "running"), + msgUpdated("m1", cost = 0.01), + updated("m1", "p1", "tool", state = "completed"), + )) + + // running is flushed when MU arrives, completed is a new batch → cannot merge + assertEquals(3, result.size) + assertEquals("running", (result[0] as ChatEventDto.PartUpdated).part.state) + assertNotNull(result[1] as? ChatEventDto.MessageUpdated) + assertEquals("completed", (result[2] as ChatEventDto.PartUpdated).part.state) + } + + // ------ helpers ------ + + private fun msgUpdated(id: String, role: String = "assistant", cost: Double? = null) = + ChatEventDto.MessageUpdated( + "ses", + MessageDto(id = id, sessionID = "ses", role = role, time = MessageTimeDto(0.0), cost = cost), + ) + + private fun statusChanged(type: String) = + ChatEventDto.SessionStatusChanged("ses", SessionStatusDto(type)) } diff --git a/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionUpdateQueueTest.kt b/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionUpdateQueueTest.kt index b59272c5758..f1076e7f09c 100644 --- a/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionUpdateQueueTest.kt +++ b/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionUpdateQueueTest.kt @@ -1,5 +1,7 @@ package ai.kilocode.client.session +import ai.kilocode.client.session.model.Tool +import ai.kilocode.client.session.model.ToolExecState import ai.kilocode.client.session.model.SessionModelEvent import ai.kilocode.client.session.model.SessionState import ai.kilocode.rpc.dto.ChatEventDto @@ -78,4 +80,101 @@ class SessionUpdateQueueTest : SessionControllerTestBase() { assertTrue(modelEvents.any { it is SessionModelEvent.StateChanged }) assertTrue(m.model.state is SessionState.Busy) } + + fun `test buffered part updates for new part collapse to one content add`() { + appRpc.state.value = ai.kilocode.rpc.dto.KiloAppStateDto(ai.kilocode.rpc.dto.KiloAppStatusDto.READY) + projectRpc.state.value = workspaceReady() + val m = controller("ses_test", flushMs = Long.MAX_VALUE) + val modelEvents = collectModelEvents(m) + flush() + modelEvents.clear() + + emit(ChatEventDto.MessageUpdated("ses_test", msg("msg1", "ses_test", "assistant"))) + modelEvents.clear() + + emit(ChatEventDto.PartUpdated("ses_test", part("prt1", "ses_test", "msg1", "tool", tool = "bash", state = "pending")), flush = false) + emit(ChatEventDto.PartUpdated("ses_test", part("prt1", "ses_test", "msg1", "tool", tool = "bash", state = "completed")), flush = false) + settle() + flush() + + assertEquals(1, modelEvents.count { it is SessionModelEvent.ContentAdded }) + assertEquals(0, modelEvents.count { it is SessionModelEvent.ContentUpdated }) + val tool = m.model.message("msg1")!!.parts["prt1"] as Tool + assertEquals(ToolExecState.COMPLETED, tool.state) + } + + fun `test buffered part updates for existing part collapse to one content update`() { + appRpc.state.value = ai.kilocode.rpc.dto.KiloAppStateDto(ai.kilocode.rpc.dto.KiloAppStatusDto.READY) + projectRpc.state.value = workspaceReady() + val m = controller("ses_test", flushMs = Long.MAX_VALUE) + val modelEvents = collectModelEvents(m) + flush() + modelEvents.clear() + + emit(ChatEventDto.MessageUpdated("ses_test", msg("msg1", "ses_test", "assistant"))) + emit(ChatEventDto.PartUpdated("ses_test", part("prt1", "ses_test", "msg1", "tool", tool = "bash", state = "pending"))) + modelEvents.clear() + + emit(ChatEventDto.PartUpdated("ses_test", part("prt1", "ses_test", "msg1", "tool", tool = "bash", state = "running")), flush = false) + emit(ChatEventDto.PartUpdated("ses_test", part("prt1", "ses_test", "msg1", "tool", tool = "bash", state = "completed", title = "Install deps")), flush = false) + settle() + flush() + + assertEquals(0, modelEvents.count { it is SessionModelEvent.ContentAdded }) + assertEquals(1, modelEvents.count { it is SessionModelEvent.ContentUpdated }) + val tool = m.model.message("msg1")!!.parts["prt1"] as Tool + assertEquals(ToolExecState.COMPLETED, tool.state) + assertEquals("Install deps", tool.title) + } + + fun `test buffered same part tool updates keep only final busy text`() { + appRpc.state.value = ai.kilocode.rpc.dto.KiloAppStateDto(ai.kilocode.rpc.dto.KiloAppStatusDto.READY) + projectRpc.state.value = workspaceReady() + val m = controller("ses_test", flushMs = Long.MAX_VALUE) + val modelEvents = collectModelEvents(m) + flush() + modelEvents.clear() + + emit(ChatEventDto.TurnOpen("ses_test")) + emit(ChatEventDto.MessageUpdated("ses_test", msg("msg1", "ses_test", "assistant"))) + modelEvents.clear() + + emit(ChatEventDto.PartUpdated("ses_test", part("prt1", "ses_test", "msg1", "tool", tool = "read", state = "running")), flush = false) + emit(ChatEventDto.PartUpdated("ses_test", part("prt1", "ses_test", "msg1", "tool", tool = "bash", state = "running")), flush = false) + settle() + flush() + + val busy = modelEvents.filterIsInstance() + .filter { it.state is SessionState.Busy } + assertEquals(1, busy.size) + val state = busy.single().state as SessionState.Busy + assertTrue(state.text.contains("commands", ignoreCase = true)) + } + + fun `test barrier prevents part update merge across turn close`() { + appRpc.state.value = ai.kilocode.rpc.dto.KiloAppStateDto(ai.kilocode.rpc.dto.KiloAppStatusDto.READY) + projectRpc.state.value = workspaceReady() + val m = controller("ses_test", flushMs = Long.MAX_VALUE) + val modelEvents = collectModelEvents(m) + flush() + modelEvents.clear() + + emit(ChatEventDto.MessageUpdated("ses_test", msg("msg1", "ses_test", "assistant"))) + emit(ChatEventDto.TurnOpen("ses_test")) + modelEvents.clear() + + emit(ChatEventDto.PartUpdated("ses_test", part("prt1", "ses_test", "msg1", "tool", tool = "bash", state = "running")), flush = false) + emit(ChatEventDto.TurnClose("ses_test", "completed"), flush = false) + emit(ChatEventDto.PartUpdated("ses_test", part("prt1", "ses_test", "msg1", "tool", tool = "bash", state = "completed")), flush = false) + settle() + flush() + + assertModelEvents(""" + ContentAdded msg1/prt1 + StateChanged Busy + StateChanged Idle + ContentUpdated msg1/prt1 + """, modelEvents) + assertEquals(SessionState.Idle, m.model.state) + } } From 1b95ef8539938b752179937a5cd0f2551924e237 Mon Sep 17 00:00:00 2001 From: kirillk Date: Thu, 23 Apr 2026 12:33:35 -0400 Subject: [PATCH 4/9] test(jetbrains): add condense toggle and parity test for condensed vs raw event delivery Add a condense: Boolean = true flag to SessionUpdateQueue and expose it as an optional param on SessionController so tests can run both modes. Default remains true, preserving existing production behaviour. Add SessionControllerTestBase.snapshot() which captures final model state (transcript, turns, diff, todos, compactionCount) for assertion. Add a large corpus parity test that feeds 49 events through a condensed and a raw controller, then asserts both reach identical final state. --- .../client/session/SessionController.kt | 3 +- .../client/session/SessionUpdateQueue.kt | 4 +- .../session/SessionControllerTestBase.kt | 34 ++++++++- .../client/session/SessionUpdateQueueTest.kt | 73 +++++++++++++++++++ 4 files changed, 110 insertions(+), 4 deletions(-) diff --git a/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionController.kt b/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionController.kt index c597e4bd7cf..bdc6303385f 100644 --- a/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionController.kt +++ b/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionController.kt @@ -57,6 +57,7 @@ class SessionController( private val cs: CoroutineScope, comp: java.awt.Component? = null, private val flushMs: Long = EVENT_FLUSH_MS, + private val condense: Boolean = true, ) : Disposable { companion object { @@ -72,7 +73,7 @@ class SessionController( private val listeners = mutableListOf() private var sessionId: String? = id private val directory: String get() = workspace.directory - private val updates = SessionUpdateQueue(parent, comp, flushMs, ::handle, id != null) { sessionId ?: "pending" } + private val updates = SessionUpdateQueue(parent, comp, flushMs, ::handle, condense, id != null) { sessionId ?: "pending" } private var partType: String? = null private var tool: String? = null diff --git a/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionUpdateQueue.kt b/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionUpdateQueue.kt index 7a7709cfb21..dc9d92cc72f 100644 --- a/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionUpdateQueue.kt +++ b/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionUpdateQueue.kt @@ -18,6 +18,7 @@ internal class SessionUpdateQueue( private val comp: Component?, private val flushMs: Long = EVENT_FLUSH_MS, private val fire: (List) -> Unit, + private val condense: Boolean = true, hold: Boolean, private val sid: () -> String, ) : Disposable { @@ -80,7 +81,7 @@ internal class SessionUpdateQueue( val before = pending.size val types = pending.groupBy { it::class.simpleName } .entries.joinToString(",") { (k, v) -> "$k:${v.size}" } - val batch = condenser.condense(pending.toList()) + val batch = if (condense) condenser.condense(pending.toList()) else pending.toList() pending.clear() last = now LOG.debug { "${ChatLogSummary.sid(sid())} flush source=$source forced=$forced pending=$before condensed=${batch.size} saved=${before - batch.size} types=$types" } @@ -98,4 +99,3 @@ internal class SessionUpdateQueue( } } - diff --git a/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionControllerTestBase.kt b/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionControllerTestBase.kt index 8426a028b7c..b59e60ac058 100644 --- a/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionControllerTestBase.kt +++ b/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionControllerTestBase.kt @@ -4,6 +4,7 @@ import ai.kilocode.client.app.KiloAppService import ai.kilocode.client.app.KiloSessionService import ai.kilocode.client.session.model.SessionModel import ai.kilocode.client.session.model.SessionModelEvent +import ai.kilocode.client.session.model.SessionState import ai.kilocode.client.testing.FakeAppRpcApi import ai.kilocode.client.testing.FakeWorkspaceRpcApi import ai.kilocode.client.testing.FakeSessionRpcApi @@ -41,6 +42,24 @@ import kotlinx.coroutines.runBlocking */ abstract class SessionControllerTestBase : BasePlatformTestCase() { + protected data class Snapshot( + val body: String, + val turns: String, + val state: SessionState, + val diff: List, + val todos: List, + val compacted: Int, + ) { + override fun toString(): String = buildString { + appendLine("state=$state") + appendLine("turns=$turns") + appendLine("diff=$diff") + appendLine("todos=$todos") + appendLine("compacted=$compacted") + append("body=\n$body") + } + } + private class Root : javax.swing.JPanel() { private var shown = true override fun isShowing(): Boolean = shown @@ -93,8 +112,12 @@ abstract class SessionControllerTestBase : BasePlatformTestCase() { protected fun controller(id: String? = null) = controller(id, Long.MAX_VALUE) protected fun controller(id: String? = null, flushMs: Long): SessionController { + return controller(id, flushMs, true) + } + + protected fun controller(id: String? = null, flushMs: Long, condense: Boolean): SessionController { val root = Root() - val m = SessionController(parent, id, sessions, workspace, app, scope, root, flushMs) + val m = SessionController(parent, id, sessions, workspace, app, scope, root, flushMs, condense) controllers.add(m) roots[m] = root return m @@ -198,6 +221,15 @@ abstract class SessionControllerTestBase : BasePlatformTestCase() { assertEquals(expected.trimIndent().trim(), events.joinToString("\n")) } + protected fun snapshot(c: SessionController) = Snapshot( + body = c.model.toString().trim(), + turns = c.model.toTurnsString().trim(), + state = c.model.state, + diff = c.model.diff.toList(), + todos = c.model.todos.toList(), + compacted = c.model.compactionCount, + ) + // ------ DTO factories ------ protected fun msg(id: String, sid: String, role: String) = MessageDto( diff --git a/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionUpdateQueueTest.kt b/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionUpdateQueueTest.kt index f1076e7f09c..8be8c905a75 100644 --- a/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionUpdateQueueTest.kt +++ b/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionUpdateQueueTest.kt @@ -5,6 +5,9 @@ import ai.kilocode.client.session.model.ToolExecState import ai.kilocode.client.session.model.SessionModelEvent import ai.kilocode.client.session.model.SessionState import ai.kilocode.rpc.dto.ChatEventDto +import ai.kilocode.rpc.dto.DiffFileDto +import ai.kilocode.rpc.dto.SessionStatusDto +import ai.kilocode.rpc.dto.TodoDto class SessionUpdateQueueTest : SessionControllerTestBase() { @@ -177,4 +180,74 @@ class SessionUpdateQueueTest : SessionControllerTestBase() { """, modelEvents) assertEquals(SessionState.Idle, m.model.state) } + + fun `test condensed and raw controller end with same final state on large corpus`() { + appRpc.state.value = ai.kilocode.rpc.dto.KiloAppStateDto(ai.kilocode.rpc.dto.KiloAppStatusDto.READY) + projectRpc.state.value = workspaceReady() + + val events = corpus() + val condensed = runCorpus(events, true) + val raw = runCorpus(events, false) + val a = snapshot(condensed) + val b = snapshot(raw) + + if (a != b) fail("condensed=\n$a\nraw=\n$b") + assertEquals(SessionState.Idle, a.state) + assertTrue(a.body.contains("assistant#msg1")) + assertTrue(a.body.contains("assistant#msg2")) + assertTrue(a.body.contains("diff: src/A.kt src/B.kt")) + assertTrue(a.body.contains("todo: [completed] ship feature")) + assertEquals(4, a.compacted) + } + + private fun corpus(): List = buildList { + add(ChatEventDto.TurnOpen("ses_test")) + add(ChatEventDto.MessageUpdated("ses_test", msg("msg1", "ses_test", "assistant"))) + add(ChatEventDto.MessageUpdated("ses_test", msg("msg1", "ses_test", "assistant").copy(cost = 0.01))) + add(ChatEventDto.MessageUpdated("ses_test", msg("msg1", "ses_test", "assistant").copy(cost = 0.02))) + add(ChatEventDto.PartUpdated("ses_test", part("tool1", "ses_test", "msg1", "tool", tool = "read", state = "running"))) + add(ChatEventDto.PartUpdated("ses_test", part("tool1", "ses_test", "msg1", "tool", tool = "read", state = "running", title = "Read files"))) + add(ChatEventDto.PartUpdated("ses_test", part("tool1", "ses_test", "msg1", "tool", tool = "read", state = "completed", title = "Read files"))) + add(ChatEventDto.PartUpdated("ses_test", part("snap1", "ses_test", "msg1", "text", text = "he"))) + repeat(8) { i -> + add(ChatEventDto.PartDelta("ses_test", "msg1", "txt1", "text", " chunk$i")) + } + add(ChatEventDto.PartUpdated("ses_test", part("snap1", "ses_test", "msg1", "text", text = "hello"))) + add(ChatEventDto.SessionStatusChanged("ses_test", SessionStatusDto("busy"))) + add(ChatEventDto.SessionStatusChanged("ses_test", SessionStatusDto("retry", message = "retrying", attempt = 2, next = 10L))) + add(ChatEventDto.SessionStatusChanged("ses_test", SessionStatusDto("offline", message = "offline", requestID = "req1"))) + add(ChatEventDto.SessionStatusChanged("ses_test", SessionStatusDto("idle"))) + add(ChatEventDto.SessionDiffChanged("ses_test", listOf(DiffFileDto("src/A.kt", 1, 0)))) + add(ChatEventDto.SessionDiffChanged("ses_test", emptyList())) + add(ChatEventDto.SessionDiffChanged("ses_test", listOf(DiffFileDto("src/A.kt", 2, 1), DiffFileDto("src/B.kt", 4, 0)))) + add(ChatEventDto.TodoUpdated("ses_test", listOf(TodoDto("draft plan", "in_progress", "high")))) + add(ChatEventDto.TodoUpdated("ses_test", listOf(TodoDto("ship feature", "completed", "high")))) + add(ChatEventDto.SessionCompacted("ses_test")) + add(ChatEventDto.MessageUpdated("ses_test", msg("msg2", "ses_test", "assistant"))) + add(ChatEventDto.MessageUpdated("ses_test", msg("msg2", "ses_test", "assistant").copy(cost = 0.02))) + add(ChatEventDto.PartUpdated("ses_test", part("tool2", "ses_test", "msg2", "tool", tool = "edit", state = "running"))) + add(ChatEventDto.PartUpdated("ses_test", part("tool2", "ses_test", "msg2", "tool", tool = "edit", state = "completed", title = "Patch file"))) + repeat(6) { i -> + add(ChatEventDto.PartDelta("ses_test", "msg2", "txt2", "text", " body$i")) + } + add(ChatEventDto.TurnClose("ses_test", "completed")) + add(ChatEventDto.TurnOpen("ses_test")) + add(ChatEventDto.MessageUpdated("ses_test", msg("msg3", "ses_test", "assistant"))) + add(ChatEventDto.MessageUpdated("ses_test", msg("msg3", "ses_test", "assistant").copy(cost = 0.03))) + add(ChatEventDto.PartUpdated("ses_test", part("tail", "ses_test", "msg3", "text", text = "tail start"))) + repeat(5) { i -> + add(ChatEventDto.PartDelta("ses_test", "msg3", "tail", "text", " extra$i")) + } + add(ChatEventDto.SessionCompacted("ses_test")) + add(ChatEventDto.SessionIdle("ses_test")) + } + + private fun runCorpus(events: List, condense: Boolean): SessionController { + val m = controller("ses_test", flushMs = Long.MAX_VALUE, condense = condense) + flush() + for (event in events) emit(event, flush = false) + settle() + flush() + return m + } } From 004f97a6cc9170e7e445b846cb4936300bfe3c1c Mon Sep 17 00:00:00 2001 From: kirillk Date: Thu, 23 Apr 2026 12:38:32 -0400 Subject: [PATCH 5/9] feat(jetbrains): read condense and flushMs from IntelliJ registry in SessionUi Register kilo.session.condense (bool, default true) and kilo.session.flushMs (int, default 150) as platform registry keys so developers can tune or disable event condensing at runtime without recompiling. --- .../main/kotlin/ai/kilocode/client/session/SessionUi.kt | 7 ++++++- .../src/main/resources/kilo.jetbrains.frontend.xml | 7 +++++++ 2 files changed, 13 insertions(+), 1 deletion(-) diff --git a/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionUi.kt b/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionUi.kt index 4bf70fa7d51..90672de6f9f 100644 --- a/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionUi.kt +++ b/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionUi.kt @@ -13,6 +13,7 @@ import ai.kilocode.client.session.ui.SessionPanel import ai.kilocode.client.session.ui.StatusPanel import com.intellij.openapi.Disposable import com.intellij.openapi.project.Project +import com.intellij.openapi.util.registry.Registry import ai.kilocode.log.ChatLogSummary import ai.kilocode.log.KiloLog import com.intellij.ui.components.JBScrollPane @@ -52,7 +53,11 @@ class SessionUi( private val LOG = KiloLog.create(SessionUi::class.java) } - private val controller = SessionController(this, null, sessions, workspace, app, cs, this) + private val controller = SessionController( + this, null, sessions, workspace, app, cs, this, + flushMs = Registry.intValue("kilo.session.flushMs", EVENT_FLUSH_MS.toInt()).toLong(), + condense = Registry.`is`("kilo.session.condense", true), + ) // ------ card switch ------ diff --git a/packages/kilo-jetbrains/frontend/src/main/resources/kilo.jetbrains.frontend.xml b/packages/kilo-jetbrains/frontend/src/main/resources/kilo.jetbrains.frontend.xml index d1815656660..cc1a61b5649 100644 --- a/packages/kilo-jetbrains/frontend/src/main/resources/kilo.jetbrains.frontend.xml +++ b/packages/kilo-jetbrains/frontend/src/main/resources/kilo.jetbrains.frontend.xml @@ -11,6 +11,13 @@ anchor="left" icon="/icons/kilo.svg" factoryClass="ai.kilocode.client.KiloToolWindowFactory"/> + + + From d11905fce42af9c2a0142d1a43d68b266b8f47f9 Mon Sep 17 00:00:00 2001 From: kirillk Date: Thu, 23 Apr 2026 12:55:33 -0400 Subject: [PATCH 6/9] feat(jetbrains): mark session registry keys as runtime-tunable Mark kilo.session.condense and kilo.session.flushMs as non-restart registry keys so queue tuning can be adjusted live while profiling session update behavior in the JetBrains client. --- .../src/main/resources/kilo.jetbrains.frontend.xml | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/packages/kilo-jetbrains/frontend/src/main/resources/kilo.jetbrains.frontend.xml b/packages/kilo-jetbrains/frontend/src/main/resources/kilo.jetbrains.frontend.xml index cc1a61b5649..5922e16ebd1 100644 --- a/packages/kilo-jetbrains/frontend/src/main/resources/kilo.jetbrains.frontend.xml +++ b/packages/kilo-jetbrains/frontend/src/main/resources/kilo.jetbrains.frontend.xml @@ -14,10 +14,14 @@ + defaultValue="true" + restartRequired="false" + overrides="false"/> + defaultValue="150" + restartRequired="false" + overrides="false"/> From 38b07b779b54e430e21a60102a7adb5ce6002f8d Mon Sep 17 00:00:00 2001 From: kirillk Date: Thu, 23 Apr 2026 13:17:40 -0400 Subject: [PATCH 7/9] fix(jetbrains): guard registry flush config and condense hidden queues Clamp non-positive kilo.session.flushMs values back to the default hardcoded cadence so invalid registry overrides cannot break session UI construction. Condense pending hidden-session events in memory without flushing them, which caps backlog growth while preserving the existing visibility gate. Add tests that verify hidden controllers still avoid model delivery until shown, while benefiting from pre-flush condensation. --- .../ai/kilocode/client/session/SessionUi.kt | 7 +++- .../client/session/SessionUpdateQueue.kt | 12 +++++- .../client/session/SessionUpdateQueueTest.kt | 42 +++++++++++++++++++ 3 files changed, 59 insertions(+), 2 deletions(-) diff --git a/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionUi.kt b/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionUi.kt index 90672de6f9f..b9b4856c4d2 100644 --- a/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionUi.kt +++ b/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionUi.kt @@ -53,9 +53,14 @@ class SessionUi( private val LOG = KiloLog.create(SessionUi::class.java) } + private val flushMs = Registry.intValue("kilo.session.flushMs", EVENT_FLUSH_MS.toInt()) + .takeIf { it > 0 } + ?.toLong() + ?: EVENT_FLUSH_MS + private val controller = SessionController( this, null, sessions, workspace, app, cs, this, - flushMs = Registry.intValue("kilo.session.flushMs", EVENT_FLUSH_MS.toInt()).toLong(), + flushMs = flushMs, condense = Registry.`is`("kilo.session.condense", true), ) diff --git a/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionUpdateQueue.kt b/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionUpdateQueue.kt index dc9d92cc72f..3496d68ba15 100644 --- a/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionUpdateQueue.kt +++ b/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionUpdateQueue.kt @@ -74,6 +74,7 @@ internal class SessionUpdateQueue( private fun flushNow(forced: Boolean, source: String) { if (hold) return + condenseHidden() if (!showing()) return if (pending.isEmpty()) return val now = System.currentTimeMillis() @@ -88,6 +89,16 @@ internal class SessionUpdateQueue( fire(batch) } + private fun condenseHidden() { + if (!condense) return + if (showing()) return + if (pending.size < 2) return + val batch = condenser.condense(pending.toList()) + if (batch.size == pending.size) return + pending.clear() + pending.addAll(batch) + } + private fun showing(): Boolean = comp?.isShowing ?: true private fun edt(block: () -> Unit) { @@ -98,4 +109,3 @@ internal class SessionUpdateQueue( app.invokeLater(block) } } - diff --git a/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionUpdateQueueTest.kt b/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionUpdateQueueTest.kt index 8be8c905a75..8b5d56e0f0b 100644 --- a/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionUpdateQueueTest.kt +++ b/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionUpdateQueueTest.kt @@ -39,6 +39,48 @@ class SessionUpdateQueueTest : SessionControllerTestBase() { assertTrue(m.model.state is SessionState.Busy) } + fun `test hidden controller condenses while hidden but does not flush`() { + appRpc.state.value = ai.kilocode.rpc.dto.KiloAppStateDto(ai.kilocode.rpc.dto.KiloAppStatusDto.READY) + projectRpc.state.value = workspaceReady() + val m = controller("ses_test", flushMs = 250L) + val modelEvents = collectModelEvents(m) + flush() + modelEvents.clear() + + hide(m) + emit(ChatEventDto.MessageUpdated("ses_test", msg("msg1", "ses_test", "assistant")), flush = false) + repeat(4) { i -> + emit(ChatEventDto.PartDelta("ses_test", "msg1", "txt1", "text", " chunk$i"), flush = false) + } + emit(ChatEventDto.PartUpdated("ses_test", part("tool1", "ses_test", "msg1", "tool", tool = "bash", state = "running")), flush = false) + emit(ChatEventDto.PartUpdated("ses_test", part("tool1", "ses_test", "msg1", "tool", tool = "bash", state = "completed", title = "Run build")), flush = false) + settle() + + assertTrue(modelEvents.isEmpty()) + assertEquals(SessionState.Idle, m.model.state) + + show(m) + settle() + flush() + + assertModelEvents(""" + MessageAdded msg1 + TurnAdded msg1 [msg1] + ContentAdded msg1/txt1 + ContentDelta msg1/txt1 + ContentAdded msg1/tool1 + """, modelEvents) + assertModel( + """ + assistant#msg1 + text#txt1: + chunk0 chunk1 chunk2 chunk3 + tool#tool1 bash [COMPLETED] Run build + """, + m, + ) + } + fun `test buffered deltas coalesce into one model delta`() { appRpc.state.value = ai.kilocode.rpc.dto.KiloAppStateDto(ai.kilocode.rpc.dto.KiloAppStatusDto.READY) projectRpc.state.value = workspaceReady() From fd9d5ca739fdd5e133486dbd78535ae82b3eb692 Mon Sep 17 00:00:00 2001 From: kirillk Date: Thu, 23 Apr 2026 13:22:53 -0400 Subject: [PATCH 8/9] fix(cli): add local jschardet module declaration for typecheck Add a minimal ambient module declaration for jschardet so the repo pre-push typecheck hook passes in worktrees where upstream dependency typings are not resolved by tsgo. --- packages/opencode/src/jschardet.d.ts | 14 ++++++++++++++ 1 file changed, 14 insertions(+) create mode 100644 packages/opencode/src/jschardet.d.ts diff --git a/packages/opencode/src/jschardet.d.ts b/packages/opencode/src/jschardet.d.ts new file mode 100644 index 00000000000..7e819af9cc1 --- /dev/null +++ b/packages/opencode/src/jschardet.d.ts @@ -0,0 +1,14 @@ +declare module "jschardet" { + export interface Result { + encoding?: string + confidence?: number + } + + export function detect(input: ArrayLike): Result + + const api: { + detect(input: ArrayLike): Result + } + + export default api +} From 0a142c5473f04ec5dbc9e286043d31e4c2429bf3 Mon Sep 17 00:00:00 2001 From: kirillk Date: Thu, 23 Apr 2026 14:55:52 -0400 Subject: [PATCH 9/9] fix(jetbrains): stop hidden sessions from loading the EDT MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit While a session panel is hidden, enqueue and timer-tick paths no longer schedule EDT runnables. Events accumulate off-EDT behind a lock, and a single visibility listener forces one forced flush when the component becomes showing again. Removes the per-enqueue condenseHidden O(n²) rebuild entirely. Closes #9437 --- .../client/session/SessionUpdateQueue.kt | 69 ++++++++++++------- .../session/SessionControllerTestBase.kt | 11 +++ .../client/session/SessionUpdateQueueTest.kt | 63 ++++++++++++++++- 3 files changed, 115 insertions(+), 28 deletions(-) diff --git a/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionUpdateQueue.kt b/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionUpdateQueue.kt index 3496d68ba15..2c02e7b3349 100644 --- a/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionUpdateQueue.kt +++ b/packages/kilo-jetbrains/frontend/src/main/kotlin/ai/kilocode/client/session/SessionUpdateQueue.kt @@ -7,9 +7,12 @@ import com.intellij.openapi.Disposable import com.intellij.openapi.application.ApplicationManager import com.intellij.openapi.util.Disposer import java.awt.Component +import java.awt.event.HierarchyEvent +import java.awt.event.HierarchyListener import java.util.concurrent.Executors import java.util.concurrent.ScheduledExecutorService import java.util.concurrent.TimeUnit +import java.util.concurrent.atomic.AtomicBoolean internal const val EVENT_FLUSH_MS = 150L @@ -29,14 +32,26 @@ internal class SessionUpdateQueue( private val app = ApplicationManager.getApplication() private val condenser = SessionQueueCondenser() private val pending = mutableListOf() + private val lock = Any() private val exec: ScheduledExecutorService? = if (flushMs == Long.MAX_VALUE) null else Executors.newSingleThreadScheduledExecutor() + private val visible = AtomicBoolean(comp?.isShowing ?: true) + private val watch = comp?.let { + HierarchyListener { event -> + if (event.changeFlags and HierarchyEvent.SHOWING_CHANGED.toLong() == 0L) return@HierarchyListener + onVisible(it.isShowing) + } + } private var last = 0L private var hold = hold init { Disposer.register(parent, this) + if (comp != null && watch != null) comp.addHierarchyListener(watch) exec?.scheduleAtFixedRate( - { requestFlush(false, "tick") }, + { + if (!visible.get()) return@scheduleAtFixedRate + requestFlush(false, "tick") + }, flushMs, flushMs, TimeUnit.MILLISECONDS, @@ -44,11 +59,13 @@ internal class SessionUpdateQueue( } fun enqueue(event: ChatEventDto) { - edt { - LOG.debug { "${ChatLogSummary.sid(sid())} enqueue pending=${pending.size + 1}" } + val size = synchronized(lock) { pending.add(event) - flushNow(false, "enqueue") + pending.size } + LOG.debug { "${ChatLogSummary.sid(sid())} enqueue pending=$size visible=${visible.get()}" } + if (!visible.get()) return + requestFlush(false, "enqueue") } fun holdFlush(hold: Boolean) { @@ -59,48 +76,48 @@ internal class SessionUpdateQueue( } fun requestFlush(forced: Boolean, source: String = "api") { + if (!forced && !visible.get()) return edt { flushNow(forced, source) } } override fun dispose() { - LOG.debug { "${ChatLogSummary.sid(sid())} dispose pending=${pending.size}" } + val size = synchronized(lock) { pending.size } + LOG.debug { "${ChatLogSummary.sid(sid())} dispose pending=$size" } exec?.shutdownNow() + if (comp != null && watch != null) comp.removeHierarchyListener(watch) if (app.isDispatchThread) { - pending.clear() + synchronized(lock) { pending.clear() } return } - app.invokeLater { pending.clear() } + app.invokeLater { synchronized(lock) { pending.clear() } } } private fun flushNow(forced: Boolean, source: String) { if (hold) return - condenseHidden() - if (!showing()) return - if (pending.isEmpty()) return + if (!visible.get()) return val now = System.currentTimeMillis() if (!forced && now - last < flushMs) return - val before = pending.size - val types = pending.groupBy { it::class.simpleName } + val batch = synchronized(lock) { + if (pending.isEmpty()) return + pending.toList().also { pending.clear() } + } + val before = batch.size + val types = batch.groupBy { it::class.simpleName } .entries.joinToString(",") { (k, v) -> "$k:${v.size}" } - val batch = if (condense) condenser.condense(pending.toList()) else pending.toList() - pending.clear() + val out = if (condense) condenser.condense(batch) else batch last = now - LOG.debug { "${ChatLogSummary.sid(sid())} flush source=$source forced=$forced pending=$before condensed=${batch.size} saved=${before - batch.size} types=$types" } - fire(batch) + LOG.debug { "${ChatLogSummary.sid(sid())} flush source=$source forced=$forced pending=$before condensed=${out.size} saved=${before - out.size} types=$types" } + fire(out) } - private fun condenseHidden() { - if (!condense) return - if (showing()) return - if (pending.size < 2) return - val batch = condenser.condense(pending.toList()) - if (batch.size == pending.size) return - pending.clear() - pending.addAll(batch) + private fun onVisible(show: Boolean) { + val prev = visible.getAndSet(show) + if (prev == show) return + LOG.debug { "${ChatLogSummary.sid(sid())} visible=$show" } + if (!show) return + requestFlush(true, "visible") } - private fun showing(): Boolean = comp?.isShowing ?: true - private fun edt(block: () -> Unit) { if (app.isDispatchThread) { block() diff --git a/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionControllerTestBase.kt b/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionControllerTestBase.kt index b59e60ac058..8e79cec7dc9 100644 --- a/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionControllerTestBase.kt +++ b/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionControllerTestBase.kt @@ -28,6 +28,7 @@ import com.intellij.openapi.application.ApplicationManager import com.intellij.openapi.util.Disposer import com.intellij.testFramework.fixtures.BasePlatformTestCase import com.intellij.util.ui.UIUtil +import java.awt.event.HierarchyEvent import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.cancel @@ -64,7 +65,17 @@ abstract class SessionControllerTestBase : BasePlatformTestCase() { private var shown = true override fun isShowing(): Boolean = shown fun showState(show: Boolean) { + val prev = shown shown = show + if (prev == show) return + val event = HierarchyEvent( + this, + HierarchyEvent.HIERARCHY_CHANGED, + this, + this.parent, + HierarchyEvent.SHOWING_CHANGED.toLong(), + ) + hierarchyListeners.forEach { it.hierarchyChanged(event) } } } diff --git a/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionUpdateQueueTest.kt b/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionUpdateQueueTest.kt index 8b5d56e0f0b..04dc2ad4f3e 100644 --- a/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionUpdateQueueTest.kt +++ b/packages/kilo-jetbrains/frontend/src/test/kotlin/ai/kilocode/client/session/SessionUpdateQueueTest.kt @@ -29,7 +29,6 @@ class SessionUpdateQueueTest : SessionControllerTestBase() { show(m) settle() - flush() assertModelEvents(""" StateChanged Busy @@ -61,7 +60,6 @@ class SessionUpdateQueueTest : SessionControllerTestBase() { show(m) settle() - flush() assertModelEvents(""" MessageAdded msg1 @@ -81,6 +79,67 @@ class SessionUpdateQueueTest : SessionControllerTestBase() { ) } + fun `test hidden cadence does not flush until shown`() { + appRpc.state.value = ai.kilocode.rpc.dto.KiloAppStateDto(ai.kilocode.rpc.dto.KiloAppStatusDto.READY) + projectRpc.state.value = workspaceReady() + val m = controller("ses_test", flushMs = 50L) + val modelEvents = collectModelEvents(m) + flush() + modelEvents.clear() + + hide(m) + emit(ChatEventDto.TurnOpen("ses_test"), flush = false) + emit(ChatEventDto.MessageUpdated("ses_test", msg("msg1", "ses_test", "assistant")), flush = false) + settle() + + assertTrue(modelEvents.isEmpty()) + assertEquals(SessionState.Idle, m.model.state) + + show(m) + settle() + + assertModelEvents(""" + StateChanged Busy + MessageAdded msg1 + TurnAdded msg1 [msg1] + """, modelEvents) + } + + fun `test hidden controller flushes on show without new event`() { + appRpc.state.value = ai.kilocode.rpc.dto.KiloAppStateDto(ai.kilocode.rpc.dto.KiloAppStatusDto.READY) + projectRpc.state.value = workspaceReady() + val m = controller("ses_test", flushMs = 250L) + val modelEvents = collectModelEvents(m) + flush() + modelEvents.clear() + + hide(m) + emit(ChatEventDto.MessageUpdated("ses_test", msg("msg1", "ses_test", "assistant")), flush = false) + emit(ChatEventDto.PartDelta("ses_test", "msg1", "txt1", "text", "hello "), flush = false) + emit(ChatEventDto.PartDelta("ses_test", "msg1", "txt1", "text", "world"), flush = false) + settle() + + assertTrue(modelEvents.isEmpty()) + + show(m) + settle() + + assertModelEvents(""" + MessageAdded msg1 + TurnAdded msg1 [msg1] + ContentAdded msg1/txt1 + ContentDelta msg1/txt1 + """, modelEvents) + assertModel( + """ + assistant#msg1 + text#txt1: + hello world + """, + m, + ) + } + fun `test buffered deltas coalesce into one model delta`() { appRpc.state.value = ai.kilocode.rpc.dto.KiloAppStateDto(ai.kilocode.rpc.dto.KiloAppStatusDto.READY) projectRpc.state.value = workspaceReady()