Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/quiet-manual-interruptions.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
"kilo-code": patch
---

Stop manually aborted turns without briefly showing an interruption warning.
11 changes: 9 additions & 2 deletions packages/kilo-vscode/src/KiloProvider.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4127,15 +4127,22 @@ export class KiloProvider implements vscode.WebviewViewProvider, TelemetryProper
this.cancelRetry(sid)
const client = this.client
if (!client) return Promise.resolve(false)
return this.aborts.stop(client, sid, this.getWorkspaceDirectory(sid))
const directory = this.getWorkspaceDirectory(sid)
const dirs = this.aborts.directories(sid, directory)
const ids = new Map(dirs.map((dir) => [dir, this.connectionService.beginExplicitAbort(sid, dir)]))
return this.aborts.stop(client, sid, directory, dirs).then((result) => {
for (const attempt of result.attempts) {
this.connectionService.finishExplicitAbort(sid, attempt.dir, ids.get(attempt.dir)!, attempt.aborted)
}
return result.complete
})
}

private async handleAbort(sessionID?: string): Promise<void> {
const sid = sessionID || this.currentSession?.id
if (!sid || !(await this.stopSession(sid))) return
this.sessionStatusMap.set(sid, "idle")
this.streams.flush(sid)
this.postMessage({ type: "sessionTurnClosed", sessionID: sid, reason: "interrupted" })
this.postMessage({ type: "sessionStatus", sessionID: sid, status: "idle" })
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1804,7 +1804,7 @@ export class AgentManagerProvider implements Disposable {
await continueInWorktree(
{
root,
getClient: () => this.connectionService.getClient(),
connection: this.connectionService,
createWorktreeOnDisk: (opts) => this.createWorktreeOnDisk(opts),
runSetupScript: (p, b, id) => this.runSetupScriptForWorktree(p, b, id),
cleanupWorktree: async (id) => {
Expand Down
18 changes: 12 additions & 6 deletions packages/kilo-vscode/src/agent-manager/continue-in-worktree.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,10 @@ import { recordForkHandoff } from "./fork-handoff"

export interface ContinueContext {
root: string
getClient: () => KiloClient
connection: {
getClient: () => KiloClient
runExplicitAbort: <T>(sessionId: string, directory: string, action: () => Promise<T>) => Promise<T>
}
createWorktreeOnDisk: (opts: { baseBranch: string; baseRef: string }) => Promise<{
worktree: { id: string }
result: CreateWorktreeResult
Expand All @@ -30,10 +33,13 @@ export type StepResult<T> = { ok: true; value: T } | { ok: false; error: string
/** Abort a running session. Best-effort — failures are logged but not fatal. */
export async function abortSession(ctx: ContinueContext, sessionId: string): Promise<void> {
try {
const client = ctx.getClient()
await client.session.abort({ sessionID: sessionId }).catch((err) => {
ctx.log("Session abort failed (may already be idle):", getErrorMessage(err))
})
await ctx.connection
.runExplicitAbort(sessionId, ctx.root, async () => {
await ctx.connection.getClient().session.abort({ sessionID: sessionId }, { throwOnError: true })
})
.catch((err) => {
ctx.log("Session abort failed (may already be idle):", getErrorMessage(err))
})
} catch (err) {
ctx.log("Client not available for abort, continuing:", getErrorMessage(err))
}
Expand Down Expand Up @@ -96,7 +102,7 @@ async function rollback(
export async function forkSession(ctx: ContinueContext, sessionId: string, dir: string): Promise<StepResult<Session>> {
let client: KiloClient
try {
client = ctx.getClient()
client = ctx.connection.getClient()
} catch (err) {
ctx.log("Client not available for session fork:", getErrorMessage(err))
return { ok: false, error: "Not connected to CLI backend" }
Expand Down
15 changes: 11 additions & 4 deletions packages/kilo-vscode/src/kilo-provider/abort.ts
Original file line number Diff line number Diff line change
Expand Up @@ -27,20 +27,27 @@ export class SessionAbort {
this.observe(sessionID, status, dir)
}

async stop(client: KiloClient, sessionID: string, fallback: string) {
const known = this.active.has(sessionID)
directories(sessionID: string, fallback: string) {
const dirs = [...(this.active.get(sessionID) ?? [])]
if (!dirs.some((dir) => sameDirectory(dir, fallback))) dirs.push(fallback)
return dirs
}

async stop(client: KiloClient, sessionID: string, fallback: string, dirs = this.directories(sessionID, fallback)) {
const known = this.active.has(sessionID)
const results = await Promise.allSettled(dirs.map((dir) => abortSession({ client, sessionID, dir })))
const failures = results.flatMap((result, index) =>
result.status === "rejected" ? [{ dir: dirs[index], error: result.reason }] : [],
)
if (failures.length > 0) {
console.error("[Kilo New] KiloProvider: Failed to abort session in one or more directories:", failures)
return false
return {
complete: false,
attempts: results.map((result, index) => ({ dir: dirs[index], aborted: result.status === "fulfilled" })),
}
}
if (known) this.active.delete(sessionID)
return known
return { complete: known, attempts: dirs.map((dir) => ({ dir, aborted: true })) }
}

dispose(dir: string) {
Expand Down
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import { describe, expect, test } from "bun:test"
import * as vscode from "vscode"
import { KiloConnectionService } from "./connection-service"
import type { SSEPayload } from "./sdk-sse-adapter"

function state(value: boolean) {
return {
Expand Down Expand Up @@ -39,6 +40,74 @@ describe("KiloConnectionService clients", () => {
})
})

describe("KiloConnectionService explicit aborts", () => {
const close = {
id: "event-close",
type: "session.turn.close",
properties: { sessionID: "session", reason: "interrupted" },
} as SSEPayload
const status = {
type: "session.status",
properties: { sessionID: "session", status: { type: "busy" } },
} as SSEPayload

test("suppresses a successful explicit abort for every subscriber", () => {
const service = new KiloConnectionService({} as any)
const raw: SSEPayload[] = []
const first: SSEPayload[] = []
const second: SSEPayload[] = []
service.onEvent((event) => raw.push(event))
service.onEventFiltered(
() => true,
(event) => first.push(event),
)
service.onEventFiltered(
() => true,
(event) => second.push(event),
)
;(service as any).broadcast(status, "/repo")
raw.length = 0
first.length = 0
second.length = 0

const id = service.beginExplicitAbort("session", "/repo")
;(service as any).broadcast(close, "/repo")
service.finishExplicitAbort("session", "/repo", id, true)

expect(first).toEqual([])
expect(second).toEqual([])
expect(raw).toEqual([close])
})

test("replays a failed explicit abort for every subscriber", () => {
const service = new KiloConnectionService({} as any)
const raw: SSEPayload[] = []
const first: SSEPayload[] = []
const second: SSEPayload[] = []
service.onEvent((event) => raw.push(event))
service.onEventFiltered(
() => true,
(event) => first.push(event),
)
service.onEventFiltered(
() => true,
(event) => second.push(event),
)
;(service as any).broadcast(status, "/repo")
raw.length = 0
first.length = 0
second.length = 0

const id = service.beginExplicitAbort("session", "/repo")
;(service as any).broadcast(close, "/repo")
service.finishExplicitAbort("session", "/repo", id, false)

expect(first).toEqual([close])
expect(second).toEqual([close])
expect(raw).toEqual([close])
})
})

describe("KiloConnectionService viewed sessions", () => {
test("keeps Agent Manager sessions when sidebar visibility changes during a flush", async () => {
const service = new KiloConnectionService({} as any)
Expand Down
60 changes: 49 additions & 11 deletions packages/kilo-vscode/src/services/cli-backend/connection-service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import { SdkSSEAdapter, type SSEPayload } from "./sdk-sse-adapter"
import type { ServerConfig } from "./types"
import { createDuplicateEventFilter, resolveEventSessionId as resolveEventSessionIdPure } from "./connection-utils"
import { SandboxPreference } from "../sandbox-preference"
import { ExplicitAbortState } from "./explicit-abort"

export type ConnectionState = "connecting" | "connected" | "disconnected" | "error"
type SSEEventListener = (event: SSEPayload, directory?: string) => void
Expand Down Expand Up @@ -96,6 +97,8 @@ export class KiloConnectionService {
private remoteService: import("../RemoteStatusService").RemoteStatusService | null = null

private readonly eventListeners: Set<SSEEventListener> = new Set()
private readonly filteredListeners = new Set<{ filter: SSEEventFilter; listener: SSEEventListener }>()
private readonly explicitAborts = new ExplicitAbortState()
private readonly stateListeners: Set<StateListener> = new Set()
private readonly notificationDismissListeners: Set<NotificationDismissListener> = new Set()
private readonly languageChangeListeners: Set<LanguageChangeListener> = new Set()
Expand Down Expand Up @@ -276,13 +279,34 @@ export class KiloConnectionService {
* Subscribe to SSE events with a filter. The filter runs for every incoming SSE event.
*/
onEventFiltered(filter: SSEEventFilter, listener: SSEEventListener): () => void {
const wrapped: SSEEventListener = (event, directory) => {
if (!filter(event, directory)) {
return
}
listener(event, directory)
const entry = { filter, listener }
this.filteredListeners.add(entry)
return () => {
this.filteredListeners.delete(entry)
}
return this.onEvent(wrapped)
}

beginExplicitAbort(sessionID: string, directory: string): number | undefined {
return this.explicitAborts.begin(sessionID, directory)
}

finishExplicitAbort(sessionID: string, directory: string, id: number | undefined, stopped: boolean): void {
for (const item of this.explicitAborts.finish(sessionID, directory, id, stopped))
this.broadcastFiltered(item.event, item.directory)
}

async runExplicitAbort<T>(sessionID: string, directory: string, action: () => Promise<T>): Promise<T> {
const id = this.beginExplicitAbort(sessionID, directory)
return action().then(
(result) => {
this.finishExplicitAbort(sessionID, directory, id, true)
return result
},
(error) => {
this.finishExplicitAbort(sessionID, directory, id, false)
throw error
},
)
}

/**
Expand All @@ -305,6 +329,7 @@ export class KiloConnectionService {
* id after external (CLI/TUI/cascade) deletes arrive via SSE.
*/
pruneSession(sessionId: string): void {
this.explicitAborts.remove(sessionId)
for (const [mid, sid] of this.messageSessionIdsByMessageId) {
if (sid === sessionId) this.messageSessionIdsByMessageId.delete(mid)
}
Expand Down Expand Up @@ -681,6 +706,8 @@ export class KiloConnectionService {
this.sseClient?.dispose()
this.serverManager.dispose()
this.eventListeners.clear()
this.filteredListeners.clear()
this.explicitAborts.clear()
this.stateListeners.clear()
this.notificationDismissListeners.clear()
this.profileChangeListeners.clear()
Expand Down Expand Up @@ -780,6 +807,7 @@ export class KiloConnectionService {
this.stopHealthPoll()
this.stopCheckin()
const sse = this.sseClient
this.explicitAborts.clear()
this.sseClient = null
sse?.disconnect()
this.client = null
Expand Down Expand Up @@ -840,11 +868,7 @@ export class KiloConnectionService {
if (this.sseClient !== sse) return
// EventV2Bridge also emits these durable compatibility envelopes after their normal live events.
if (duplicateEvent(event)) return
this.handlePermissionEvent(event, directory)
this.handleQuestionEvent(event, directory)
for (const listener of this.eventListeners) {
listener(event, directory)
}
this.broadcast(event, directory)
})

sse.onError((error) => {
Expand Down Expand Up @@ -890,6 +914,20 @@ export class KiloConnectionService {
this.startHealthPoll(config.baseUrl, config.password)
}

private broadcast(event: SSEPayload, directory?: string): void {
this.handlePermissionEvent(event, directory)
this.handleQuestionEvent(event, directory)
for (const listener of this.eventListeners) listener(event, directory)
if (!this.explicitAborts.event(event, directory)) return
this.broadcastFiltered(event, directory)
}

private broadcastFiltered(event: SSEPayload, directory?: string): void {
for (const entry of this.filteredListeners) {
if (entry.filter(event, directory)) entry.listener(event, directory)
}
}

private startCheckin(): void {
this.stopCheckin()
this.checkinTimer = setInterval(() => this.flushViewed(), 60_000)
Expand Down
Loading
Loading