fix(acp): correlate idle events with current turn

This commit is contained in:
Johnny Eric Amancio
2026-08-19 16:34:37 +02:00
parent 860f5d9e68
commit 058f2808da
3 changed files with 145 additions and 25 deletions
+67 -13
View File
@@ -1,6 +1,7 @@
import type { AgentSideConnection } from "@agentclientprotocol/sdk"
import type {
Event,
EventMessageUpdated, // kilocode_change
EventMessagePartDelta,
EventMessagePartUpdated,
KiloClient,
@@ -41,8 +42,7 @@ export class Subscription {
private readonly shellSnapshots = new Map<string, string>()
private readonly toolStarts = new Set<string>()
private readonly connectionWaiters = new Set<() => void>()
private readonly idleWaiters = new Map<string, Set<ReturnType<typeof signal>>>()
private readonly idleCounters = new Map<string, number>() // kilocode_change
private readonly idleWaiters = new Map<string, Set<ReturnType<typeof turn>>>() // kilocode_change
private readonly permission: ACPPermission.Handler
private connected = false
private started = false
@@ -74,9 +74,8 @@ export class Subscription {
async runUntilIdle<A>(sessionId: string, request: () => Promise<A>) {
await this.waitUntilConnected()
// kilocode_change start - correlate idle waiter with per-session generation count
const start = this.idleCounters.get(sessionId) ?? 0
const waiter = signal()
// kilocode_change start - correlate idle waiter with the request's response message
const waiter = turn()
const waiters = this.idleWaiters.get(sessionId) ?? new Set()
waiters.add(waiter)
this.idleWaiters.set(sessionId, waiters)
@@ -84,7 +83,9 @@ export class Subscription {
try {
void waiter.promise.catch(() => {})
const response = await request()
if (this.connected && (this.idleCounters.get(sessionId) ?? 0) === start) {
const id = (response as { data?: { info?: { id?: string } } }).data?.info?.id
waiter.target(id ?? (await this.latest(sessionId)))
if (this.connected) {
let timer: ReturnType<typeof setTimeout> | undefined
try {
await Promise.race([
@@ -100,10 +101,7 @@ export class Subscription {
return response
} finally {
waiters.delete(waiter)
if (waiters.size === 0) {
this.idleWaiters.delete(sessionId)
this.idleCounters.delete(sessionId) // kilocode_change
}
if (waiters.size === 0) this.idleWaiters.delete(sessionId)
}
// kilocode_change end
}
@@ -116,6 +114,9 @@ export class Subscription {
case "permission.asked":
this.permission.handle(event)
return
case "message.updated":
this.message(event) // kilocode_change - correlate idle with this response message
return
case "message.part.updated":
return this.handlePartUpdated(event)
case "message.part.delta":
@@ -212,11 +213,30 @@ export class Subscription {
private idle(sessionId: string) {
const waiters = this.idleWaiters.get(sessionId)
if (!waiters) return
this.idleCounters.set(sessionId, (this.idleCounters.get(sessionId) ?? 0) + 1) // kilocode_change
this.idleWaiters.delete(sessionId)
for (const waiter of waiters) waiter.resolve()
for (const waiter of waiters) waiter.idle() // kilocode_change
}
// kilocode_change start
private async latest(sessionId: string) {
const session = await Effect.runPromise(this.input.session.tryGet(sessionId))
if (!session) throw new Error(`Missing ACP session: ${sessionId}`)
const response = await this.input.sdk.session.messages(
{ sessionID: sessionId, directory: session.cwd, limit: 1 },
{ throwOnError: true },
)
const message = response.data.at(-1)
if (!message) throw new Error(`Missing ACP response message: ${sessionId}`)
return message.info.id
}
private message(event: EventMessageUpdated) {
const sessionId = event.properties.sessionID
const waiters = this.idleWaiters.get(sessionId)
if (!waiters) return
for (const waiter of waiters) waiter.message(event.properties.info.id)
}
// kilocode_change end
private async handlePartUpdated(event: EventMessagePartUpdated) {
const part = event.properties.part
const sessionId = part.sessionID || event.properties.sessionID
@@ -447,4 +467,38 @@ function signal() {
}
}
// kilocode_change start
function turn() {
const state = {
seq: 0,
idle: 0,
target: undefined as string | undefined,
seen: new Map<string, number>(),
}
const done = signal()
const check = () => {
if (!state.target) return
const seen = state.seen.get(state.target)
if (seen === undefined || state.idle <= seen) return
done.resolve()
}
return {
promise: done.promise,
reject: done.reject,
target(id: string) {
state.target = id
check()
},
message(id: string) {
state.seen.set(id, ++state.seq)
check()
},
idle() {
state.idle = ++state.seq
check()
},
}
}
// kilocode_change end
export * as ACPEvent from "./event"
+47
View File
@@ -319,6 +319,53 @@ async function createKnownSession(
}
describe("acp event routing", () => {
// kilocode_change start
it("waits for the current turn's idle after receiving a stale idle", async () => {
const harness = createHarness()
const called = Promise.withResolvers<void>()
const response = Promise.withResolvers<{ data: { info: { id: string } } }>()
const state = { done: false }
const idle = {
id: "evt_idle",
type: "session.status",
properties: { sessionID: "ses_a", status: { type: "idle" } },
} as Event
const message = {
id: "evt_current",
type: "message.updated",
properties: { sessionID: "ses_a", info: { id: "msg_current" } },
} as Event
harness.subscription.start()
try {
await pollUntil(() => harness.calls.eventSubscribe === 1, "event stream did not connect")
await Bun.sleep(0)
const result = harness.subscription
.runUntilIdle("ses_a", () => {
called.resolve()
return response.promise
})
.then(() => {
state.done = true
})
await called.promise
await harness.subscription.handle(idle)
response.resolve({ data: { info: { id: "msg_current" } } })
await Bun.sleep(0)
await harness.subscription.handle(message)
await Bun.sleep(0)
expect(state.done).toBe(false)
await harness.subscription.handle(idle)
await result
expect(state.done).toBe(true)
} finally {
harness.subscription.stop()
}
})
// kilocode_change end
it("routes message.part.delta by sessionID without cross-session pollution", async () => {
const harness = createHarness()
await createKnownSession(harness.session, "ses_a", { messageId: "msg_a", partId: "part_a", partType: "text" })
@@ -198,6 +198,7 @@ describe("ACP service sessions", () => {
sessionUpdate?: (update: SessionNotification) => Promise<void>
},
) => {
const history = [...messages] // kilocode_change
const updates: SessionNotification[] = []
const mcpAdds: string[] = []
const aborts: string[] = []
@@ -248,7 +249,7 @@ describe("ACP service sessions", () => {
Promise.resolve({
data: input.directory ? sessions.filter((session) => session.directory === input.directory) : sessions,
}),
messages: () => Promise.resolve({ data: messages }),
messages: () => Promise.resolve({ data: history }), // kilocode_change
prompt: async (input: { sessionID: string }) => {
const response = await (options?.prompt?.(input) ??
Promise.resolve({
@@ -262,27 +263,31 @@ describe("ACP service sessions", () => {
},
}))
prompts.push(input)
events.push(updated(input.sessionID, response.data.info)) // kilocode_change
events.push(idleEvent(input.sessionID))
return response
},
command: (input: { sessionID: string }) => {
commands.push(input)
// kilocode_change start - model the response message that precedes idle
const info = assistantInfo({ input: 3, output: 4, reasoning: 0, cache: { read: 0, write: 0 } })
events.push(updated(input.sessionID, info))
events.push(idleEvent(input.sessionID))
return Promise.resolve({
data: {
info: assistantInfo({
input: 3,
output: 4,
reasoning: 0,
cache: { read: 0, write: 0 },
}),
},
})
return Promise.resolve({ data: { info } })
// kilocode_change end
},
summarize: (input: { sessionID: string }) => {
summarizes.push(input)
// kilocode_change start - model the generated summary message that precedes idle
const info = {
summary: true,
...assistantInfo({ input: 1, output: 1, reasoning: 0, cache: { read: 0, write: 0 } }),
}
history.push({ info, parts: [] })
events.push(updated(input.sessionID, info))
events.push(idleEvent(input.sessionID))
return Promise.resolve({ data: true })
// kilocode_change end
},
abort:
options?.abort ??
@@ -1341,11 +1346,15 @@ describe("ACP service sessions", () => {
})
})
// kilocode_change start - include the response identity used by the idle barrier
function assistantInfo(
tokens: UsageService.AssistantTokenCost["tokens"],
error?: AssistantMessage["error"],
): UsageService.AssistantMessage & Pick<AssistantMessage, "error"> {
): UsageService.AssistantMessage & Pick<AssistantMessage, "id" | "sessionID" | "error"> {
return {
id: "msg_assistant",
sessionID: "ses_new",
// kilocode_change end
role: "assistant",
providerID: "test",
modelID: "test-model",
@@ -1355,6 +1364,16 @@ function assistantInfo(
}
}
// kilocode_change start
function updated(sessionID: string, info: ReturnType<typeof assistantInfo>): Event {
return {
id: `evt_${info.id}`,
type: "message.updated",
properties: { sessionID, info: info as AssistantMessage },
}
}
// kilocode_change end
function categories(result: NewSessionResponse | LoadSessionResponse) {
return result.configOptions?.map((option) => option.category) ?? []
}