From 058f2808daed858dafd0046fa2eb0e17d4b1abac Mon Sep 17 00:00:00 2001 From: Johnny Eric Amancio Date: Wed, 19 Aug 2026 16:34:37 +0200 Subject: [PATCH] fix(acp): correlate idle events with current turn --- packages/opencode/src/acp/event.ts | 80 ++++++++++++++++--- packages/opencode/test/acp/event.test.ts | 47 +++++++++++ .../opencode/test/acp/service-session.test.ts | 43 +++++++--- 3 files changed, 145 insertions(+), 25 deletions(-) diff --git a/packages/opencode/src/acp/event.ts b/packages/opencode/src/acp/event.ts index 4d7ca695e76..cd7b1378cf2 100644 --- a/packages/opencode/src/acp/event.ts +++ b/packages/opencode/src/acp/event.ts @@ -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() private readonly toolStarts = new Set() private readonly connectionWaiters = new Set<() => void>() - private readonly idleWaiters = new Map>>() - private readonly idleCounters = new Map() // kilocode_change + private readonly idleWaiters = new Map>>() // kilocode_change private readonly permission: ACPPermission.Handler private connected = false private started = false @@ -74,9 +74,8 @@ export class Subscription { async runUntilIdle(sessionId: string, request: () => Promise) { 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 | 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(), + } + 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" diff --git a/packages/opencode/test/acp/event.test.ts b/packages/opencode/test/acp/event.test.ts index e860b0c57f6..0d6b5f954c4 100644 --- a/packages/opencode/test/acp/event.test.ts +++ b/packages/opencode/test/acp/event.test.ts @@ -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() + 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" }) diff --git a/packages/opencode/test/acp/service-session.test.ts b/packages/opencode/test/acp/service-session.test.ts index 46fe001503b..28187f38d50 100644 --- a/packages/opencode/test/acp/service-session.test.ts +++ b/packages/opencode/test/acp/service-session.test.ts @@ -198,6 +198,7 @@ describe("ACP service sessions", () => { sessionUpdate?: (update: SessionNotification) => Promise }, ) => { + 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 { +): UsageService.AssistantMessage & Pick { 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): 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) ?? [] }