Merge pull request #9396 from Kilo-Org/legend-trombone

feat(jetbrains): buffer and condense session updates before model delivery
This commit is contained in:
Kirill Kalishev
2026-04-23 15:15:28 -04:00
committed by GitHub
11 changed files with 1124 additions and 15 deletions
@@ -55,6 +55,9 @@ 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,
private val condense: Boolean = true,
) : Disposable {
companion object {
@@ -70,6 +73,7 @@ class SessionController(
private val listeners = mutableListOf<SessionControllerListener>()
private var sessionId: String? = id
private val directory: String get() = workspace.directory
private val updates = SessionUpdateQueue(parent, comp, flushMs, ::handle, condense, id != null) { sessionId ?: "pending" }
private var partType: String? = null
private var tool: String? = null
@@ -82,6 +86,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 +253,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 +279,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 +302,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 +468,10 @@ class SessionController(
}
}
private fun handle(events: List<ChatEventDto>) {
for (event in events) handle(event)
}
private fun showMessages() {
if (!model.showMessages) {
model.showMessages = true
@@ -500,6 +513,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()
@@ -0,0 +1,135 @@
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 same-key snapshot and text-delta events.
*
* ## Algorithm
*
* 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. 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
* - `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 {
fun condense(events: List<ChatEventDto>): List<ChatEventDto> {
if (events.size < 2) return events
val out = mutableListOf<ChatEventDto>()
val deltas = LinkedHashMap<String, ChatEventDto.PartDelta>()
val parts = LinkedHashMap<String, ChatEventDto.PartUpdated>()
val states = LinkedHashMap<String, ChatEventDto>()
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) {
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)
}
}
}
drain()
return out
}
private fun ChatEventDto.PartDelta.key(): String? {
if (field != "text") return null
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)
}
@@ -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,16 @@ class SessionUi(
private val LOG = KiloLog.create(SessionUi::class.java)
}
private val controller = SessionController(this, null, sessions, workspace, app, cs)
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 = flushMs,
condense = Registry.`is`("kilo.session.condense", true),
)
// ------ card switch ------
@@ -0,0 +1,128 @@
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
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
internal class SessionUpdateQueue(
parent: Disposable,
private val comp: Component?,
private val flushMs: Long = EVENT_FLUSH_MS,
private val fire: (List<ChatEventDto>) -> Unit,
private val condense: Boolean = true,
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<ChatEventDto>()
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(
{
if (!visible.get()) return@scheduleAtFixedRate
requestFlush(false, "tick")
},
flushMs,
flushMs,
TimeUnit.MILLISECONDS,
)
}
fun enqueue(event: ChatEventDto) {
val size = synchronized(lock) {
pending.add(event)
pending.size
}
LOG.debug { "${ChatLogSummary.sid(sid())} enqueue pending=$size visible=${visible.get()}" }
if (!visible.get()) return
requestFlush(false, "enqueue")
}
fun holdFlush(hold: Boolean) {
edt {
LOG.debug { "${ChatLogSummary.sid(sid())} hold=$hold" }
this.hold = hold
}
}
fun requestFlush(forced: Boolean, source: String = "api") {
if (!forced && !visible.get()) return
edt { flushNow(forced, source) }
}
override fun dispose() {
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) {
synchronized(lock) { pending.clear() }
return
}
app.invokeLater { synchronized(lock) { pending.clear() } }
}
private fun flushNow(forced: Boolean, source: String) {
if (hold) return
if (!visible.get()) return
val now = System.currentTimeMillis()
if (!forced && now - last < flushMs) return
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 out = if (condense) condenser.condense(batch) else batch
last = now
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 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 edt(block: () -> Unit) {
if (app.isDispatchThread) {
block()
return
}
app.invokeLater(block)
}
}
@@ -11,6 +11,17 @@
anchor="left"
icon="/icons/kilo.svg"
factoryClass="ai.kilocode.client.KiloToolWindowFactory"/>
<registryKey key="kilo.session.condense"
description="Enable event condensing in the session update queue (merges redundant snapshots before model delivery)."
defaultValue="true"
restartRequired="false"
overrides="false"/>
<registryKey key="kilo.session.flushMs"
description="Session update queue flush cadence in milliseconds."
defaultValue="150"
restartRequired="false"
overrides="false"/>
</extensions>
<actions>
@@ -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`() {
@@ -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
@@ -27,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
@@ -41,6 +43,45 @@ import kotlinx.coroutines.runBlocking
*/
abstract class SessionControllerTestBase : BasePlatformTestCase() {
protected data class Snapshot(
val body: String,
val turns: String,
val state: SessionState,
val diff: List<ai.kilocode.rpc.dto.DiffFileDto>,
val todos: List<ai.kilocode.rpc.dto.TodoDto>,
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
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) }
}
}
private val controllers = mutableListOf<SessionController>()
private val roots = mutableMapOf<SessionController, Root>()
protected lateinit var rpc: FakeSessionRpcApi
protected lateinit var appRpc: FakeAppRpcApi
protected lateinit var projectRpc: FakeWorkspaceRpcApi
@@ -79,8 +120,27 @@ 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 {
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, condense)
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 +169,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,13 +223,24 @@ abstract class SessionControllerTestBase : BasePlatformTestCase() {
}
protected fun assertControllerEvents(expected: String, events: List<SessionControllerEvent>) {
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<SessionModelEvent>) {
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(
@@ -177,6 +257,8 @@ abstract class SessionControllerTestBase : BasePlatformTestCase() {
type: String,
text: String? = null,
tool: String? = null,
state: String? = null,
title: String? = null,
) = PartDto(
id = id,
sessionID = sid,
@@ -184,6 +266,8 @@ abstract class SessionControllerTestBase : BasePlatformTestCase() {
type = type,
text = text,
tool = tool,
state = state,
title = title,
)
protected fun workspaceReady(
@@ -0,0 +1,351 @@
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() {
private val condenser = SessionQueueCondenser()
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)
fun `test empty list returns empty`() {
assertEquals(emptyList<ChatEventDto>(), 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"),
))
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))
}
@@ -0,0 +1,354 @@
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
import ai.kilocode.rpc.dto.DiffFileDto
import ai.kilocode.rpc.dto.SessionStatusDto
import ai.kilocode.rpc.dto.TodoDto
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()
assertModelEvents("""
StateChanged Busy
MessageAdded msg1
TurnAdded msg1 [msg1]
""", modelEvents)
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()
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 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()
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<SessionModelEvent.ContentDelta>()
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)
}
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<SessionModelEvent.StateChanged>()
.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)
}
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<ChatEventDto> = 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<ChatEventDto>, 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
}
}
@@ -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())
}
}
+14
View File
@@ -0,0 +1,14 @@
declare module "jschardet" {
export interface Result {
encoding?: string
confidence?: number
}
export function detect(input: ArrayLike<number>): Result
const api: {
detect(input: ArrayLike<number>): Result
}
export default api
}