diff --git a/.changeset/calm-agent-manager-streams.md b/.changeset/calm-agent-manager-streams.md new file mode 100644 index 00000000000..2d16ef00e4d --- /dev/null +++ b/.changeset/calm-agent-manager-streams.md @@ -0,0 +1,5 @@ +--- +"kilo-code": patch +--- + +Reduce duplicate event processing across VS Code when multiple sessions run concurrently. diff --git a/packages/kilo-vscode/src/services/cli-backend/connection-service.ts b/packages/kilo-vscode/src/services/cli-backend/connection-service.ts index d202eeedba9..1b0e8679d38 100644 --- a/packages/kilo-vscode/src/services/cli-backend/connection-service.ts +++ b/packages/kilo-vscode/src/services/cli-backend/connection-service.ts @@ -3,7 +3,7 @@ import { ServerManager } from "./server-manager" import { createKiloClient, type KiloClient } from "@kilocode/sdk/v2/client" import { SdkSSEAdapter, type SSEPayload } from "./sdk-sse-adapter" import type { ServerConfig } from "./types" -import { resolveEventSessionId as resolveEventSessionIdPure } from "./connection-utils" +import { createDuplicateEventFilter, resolveEventSessionId as resolveEventSessionIdPure } from "./connection-utils" import { SandboxPreference } from "../sandbox-preference" export type ConnectionState = "connecting" | "connected" | "disconnected" | "error" @@ -96,6 +96,7 @@ export class KiloConnectionService { private remoteService: import("../RemoteStatusService").RemoteStatusService | null = null private readonly eventListeners: Set = new Set() + private readonly duplicateEvent = createDuplicateEventFilter() private readonly stateListeners: Set = new Set() private readonly notificationDismissListeners: Set = new Set() private readonly languageChangeListeners: Set = new Set() @@ -837,6 +838,8 @@ export class KiloConnectionService { // Wire SSE events → broadcast to all registered listeners sse.onEvent((event, directory) => { if (this.sseClient !== sse) return + // EventV2Bridge also emits these durable compatibility envelopes after their normal live events. + if (this.duplicateEvent(event)) return this.handlePermissionEvent(event, directory) this.handleQuestionEvent(event, directory) for (const listener of this.eventListeners) { diff --git a/packages/kilo-vscode/src/services/cli-backend/connection-utils.ts b/packages/kilo-vscode/src/services/cli-backend/connection-utils.ts index f41c9a8b16f..271eec41148 100644 --- a/packages/kilo-vscode/src/services/cli-backend/connection-utils.ts +++ b/packages/kilo-vscode/src/services/cli-backend/connection-utils.ts @@ -4,6 +4,33 @@ export type { SSEPayload } from "./sdk-sse-adapter" type SyncPayload = Extract type TransientPayload = Exclude +const duplicateSyncEvents = new Set([ + "message.updated.1", + "message.removed.1", + "message.part.updated.1", + "message.part.removed.1", + "session.created.1", + "session.deleted.1", +]) + +const duplicateLiveEvents = new Set([...duplicateSyncEvents].map((name) => name.slice(0, -2))) +const DUPLICATE_EVENT_LIMIT = 1024 + +export function createDuplicateEventFilter() { + const seen = new Set() + return (event: SSEPayload): boolean => { + if (event.type === "sync") { + return duplicateSyncEvents.has(event.name) && seen.delete(event.id) + } + + if (duplicateLiveEvents.has(event.type)) { + seen.add(event.id) + if (seen.size > DUPLICATE_EVENT_LIMIT) seen.delete(seen.values().next().value!) + } + return false + } +} + /** * Pure session ID resolution for SSE events. * The lookupMessageSessionId callback remains part of the public resolver contract for diff --git a/packages/kilo-vscode/tests/unit/connection-utils.test.ts b/packages/kilo-vscode/tests/unit/connection-utils.test.ts index bd9ab804353..c1865616e7d 100644 --- a/packages/kilo-vscode/tests/unit/connection-utils.test.ts +++ b/packages/kilo-vscode/tests/unit/connection-utils.test.ts @@ -1,5 +1,5 @@ import { describe, expect, it } from "bun:test" -import { resolveEventSessionId } from "../../src/services/cli-backend/connection-utils" +import { createDuplicateEventFilter, resolveEventSessionId } from "../../src/services/cli-backend/connection-utils" import type { SSEPayload as Payload } from "../../src/services/cli-backend/sdk-sse-adapter" const noLookup = (_: string) => undefined @@ -172,3 +172,65 @@ describe("resolveEventSessionId", () => { expect(resolveEventSessionId(event, noLookup)).toBeUndefined() }) }) + +describe("isDuplicateSyncEvent", () => { + it("drops a compatibility envelope only after its live event", () => { + const filter = createDuplicateEventFilter() + const live = { + id: "e13", + type: "message.part.updated", + properties: { sessionID: "s6", part, delta: "x" }, + } satisfies Payload + expect(filter(live)).toBe(false) + expect( + filter( + sync({ + type: "sync", + name: "message.part.updated.1", + id: "e13", + seq: 5, + aggregateID: "s6", + data: { sessionID: "s6", part, time: 0 }, + }), + ), + ).toBe(true) + }) + + it("keeps replay-only compatibility envelopes", () => { + const filter = createDuplicateEventFilter() + expect( + filter( + sync({ + type: "sync", + name: "session.next.model.switched.1", + id: "e14", + seq: 6, + aggregateID: "s6", + data: { sessionID: "s6", messageID: "m1", model: { id: "test", providerID: "kilo" } }, + }), + ), + ).toBe(false) + }) + + it("keeps session updates because the provider consumes their sync metadata", () => { + const filter = createDuplicateEventFilter() + const live = { + id: "e15", + type: "session.updated", + properties: { sessionID: "s6", info: { id: "s6", time: { created: 0, updated: 1 } } }, + } as Payload + expect(filter(live)).toBe(false) + expect( + filter( + sync({ + type: "sync", + name: "session.updated.1", + id: "e15", + seq: 7, + aggregateID: "s6", + data: { sessionID: "s6", info: { id: "s6", time: { created: 0, updated: 1 } } }, + }), + ), + ).toBe(false) + }) +})