mirror of
https://github.com/Kilo-Org/kilocode.git
synced 2026-09-24 16:02:55 +08:00
fix(cli): address remote relay review feedback — race, leak, and stuck-state fixes
Fix .finally() clearing a newer enabling promise on disable→re-enable, grandchild entries leaking in the children map on unsubscribe, and permanent WS close codes leaving remote stuck enabled forever.
This commit is contained in:
@@ -268,6 +268,7 @@ export namespace KiloSessions {
|
||||
// the instance-scoped subscription map via Instance.state().
|
||||
void Instance.provide({ directory, fn: () => sender.handle(msg) })
|
||||
},
|
||||
onClose: () => disableRemote(),
|
||||
})
|
||||
|
||||
const sender = RemoteSender.create({
|
||||
@@ -289,7 +290,7 @@ export namespace KiloSessions {
|
||||
remote = { conn, sender, heartbeat }
|
||||
log.info("remote connection enabled")
|
||||
})().finally(() => {
|
||||
enabling = undefined
|
||||
if (remoteSeq === seq) enabling = undefined
|
||||
})
|
||||
|
||||
return enabling
|
||||
|
||||
@@ -93,32 +93,34 @@ export namespace RemoteSender {
|
||||
// Replay pending questions/permissions so a newly-subscribed web client
|
||||
// sees state that was asked before it connected — analogous to the Cloud
|
||||
// Agent's `connected` event carrying pending question/permission fields.
|
||||
async function replay(sessionId: string) {
|
||||
const [questions, permissions] = await Promise.all([Question.list(), PermissionNext.list()])
|
||||
for (const q of questions) {
|
||||
if (q.sessionID !== sessionId) continue
|
||||
options.conn.send({
|
||||
type: "event",
|
||||
sessionId,
|
||||
event: "question.asked",
|
||||
data: q,
|
||||
})
|
||||
}
|
||||
for (const p of permissions) {
|
||||
if (p.sessionID !== sessionId) continue
|
||||
options.conn.send({
|
||||
type: "event",
|
||||
sessionId,
|
||||
event: "permission.asked",
|
||||
data: p,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
async function backfillPendingState(sessionId: string) {
|
||||
const provide = options.provide ?? Instance.provide
|
||||
try {
|
||||
await provide({
|
||||
directory: options.directory,
|
||||
fn: async () => {
|
||||
const [questions, permissions] = await Promise.all([Question.list(), PermissionNext.list()])
|
||||
for (const q of questions) {
|
||||
if (q.sessionID !== sessionId) continue
|
||||
options.conn.send({
|
||||
type: "event",
|
||||
sessionId,
|
||||
event: "question.asked",
|
||||
data: q,
|
||||
})
|
||||
}
|
||||
for (const p of permissions) {
|
||||
if (p.sessionID !== sessionId) continue
|
||||
options.conn.send({
|
||||
type: "event",
|
||||
sessionId,
|
||||
event: "permission.asked",
|
||||
data: p,
|
||||
})
|
||||
}
|
||||
},
|
||||
fn: () => replay(sessionId),
|
||||
})
|
||||
} catch (e) {
|
||||
options.log.error("backfill pending state failed", { sessionId, error: String(e) })
|
||||
@@ -137,6 +139,7 @@ export namespace RemoteSender {
|
||||
event: "session.created",
|
||||
data: { info: child },
|
||||
})
|
||||
await replay(child.id)
|
||||
await discoverChildren(child.id)
|
||||
}
|
||||
}
|
||||
@@ -287,8 +290,15 @@ export namespace RemoteSender {
|
||||
}
|
||||
if (msg.type === "unsubscribe") {
|
||||
sessions.delete(msg.sessionId)
|
||||
for (const [child, parent] of children) {
|
||||
if (parent === msg.sessionId) children.delete(child)
|
||||
const queue = [msg.sessionId]
|
||||
while (queue.length) {
|
||||
const id = queue.pop()!
|
||||
for (const [child, parent] of children) {
|
||||
if (parent === id) {
|
||||
children.delete(child)
|
||||
queue.push(child)
|
||||
}
|
||||
}
|
||||
}
|
||||
if (sessions.size === 0 && unsub) {
|
||||
unsub()
|
||||
|
||||
@@ -16,6 +16,8 @@ export namespace RemoteWS {
|
||||
heartbeat?: number
|
||||
/** Wraps callbacks that need to run in a specific async context (e.g. Instance.provide) */
|
||||
withContext?: <R>(fn: () => R) => Promise<R> | R
|
||||
/** Called when the server permanently closes the connection (e.g. auth failure, conflict) */
|
||||
onClose?: (code: number, reason: string) => void
|
||||
}
|
||||
|
||||
export type Connection = {
|
||||
@@ -102,6 +104,7 @@ export namespace RemoteWS {
|
||||
code: event.code,
|
||||
reason: event.reason,
|
||||
})
|
||||
options.onClose?.(event.code, event.reason)
|
||||
return
|
||||
}
|
||||
schedule()
|
||||
|
||||
@@ -61,9 +61,9 @@ export function registerKiloCommands(useSDK: () => UseSDK) {
|
||||
} else {
|
||||
const result = await sdk.client.remote.enable()
|
||||
if (result.error) {
|
||||
dialog.replace(() => (
|
||||
<DialogAlert title="Not authorized" message="Run `kilo auth login` to enable remote." />
|
||||
))
|
||||
const err = result.error as { error?: string }
|
||||
const msg = err?.error ?? "Failed to enable remote."
|
||||
dialog.replace(() => <DialogAlert title="Error" message={msg} />)
|
||||
return
|
||||
}
|
||||
toast.show({ message: "Remote enabled", variant: "success" })
|
||||
|
||||
@@ -676,6 +676,48 @@ describe("RemoteSender", () => {
|
||||
expect(sent.filter((m: any) => m.event === "message.updated")).toHaveLength(0)
|
||||
})
|
||||
|
||||
test("unsubscribe parent cleans up grandchild tracking", () => {
|
||||
const { conn, sent } = fakeConn()
|
||||
const bus = fakeBus()
|
||||
const sender = RemoteSender.create({
|
||||
conn,
|
||||
directory: "/tmp/test",
|
||||
log: nolog,
|
||||
subscribe: bus.subscribe,
|
||||
})
|
||||
|
||||
sender.handle({ type: "subscribe", sessionId: "ses_root" })
|
||||
|
||||
bus.fire({
|
||||
type: "session.created",
|
||||
properties: { info: { id: "ses_child", parentID: "ses_root", title: "child" }, sessionID: "ses_child" },
|
||||
})
|
||||
bus.fire({
|
||||
type: "session.created",
|
||||
properties: {
|
||||
info: { id: "ses_grandchild", parentID: "ses_child", title: "grandchild" },
|
||||
sessionID: "ses_grandchild",
|
||||
},
|
||||
})
|
||||
|
||||
sender.handle({ type: "unsubscribe", sessionId: "ses_root" })
|
||||
sender.handle({ type: "subscribe", sessionId: "ses_keep" })
|
||||
|
||||
// Clear events from subscribe/session.created
|
||||
sent.length = 0
|
||||
|
||||
bus.fire({
|
||||
type: "message.updated",
|
||||
properties: { sessionID: "ses_child", text: "after unsub" },
|
||||
})
|
||||
bus.fire({
|
||||
type: "message.updated",
|
||||
properties: { sessionID: "ses_grandchild", text: "after unsub" },
|
||||
})
|
||||
|
||||
expect(sent).toHaveLength(0)
|
||||
})
|
||||
|
||||
test("root session events do not include parentSessionId", () => {
|
||||
const { conn, sent } = fakeConn()
|
||||
const bus = fakeBus()
|
||||
@@ -816,7 +858,14 @@ describe("RemoteSender", () => {
|
||||
metadata: {},
|
||||
always: [],
|
||||
} as any,
|
||||
{ id: "permission_2", sessionID: "ses_other", permission: "file.read", patterns: ["*"], metadata: {}, always: [] } as any,
|
||||
{
|
||||
id: "permission_2",
|
||||
sessionID: "ses_other",
|
||||
permission: "file.read",
|
||||
patterns: ["*"],
|
||||
metadata: {},
|
||||
always: [],
|
||||
} as any,
|
||||
])
|
||||
|
||||
const sender = RemoteSender.create({
|
||||
@@ -851,11 +900,16 @@ describe("RemoteSender", () => {
|
||||
const { conn, sent } = fakeConn()
|
||||
const bus = fakeBus()
|
||||
|
||||
spyOn(Question, "list").mockResolvedValue([
|
||||
{ id: "question_1", sessionID: "ses_other", questions: [] } as any,
|
||||
])
|
||||
spyOn(Question, "list").mockResolvedValue([{ id: "question_1", sessionID: "ses_other", questions: [] } as any])
|
||||
spyOn(PermissionNext, "list").mockResolvedValue([
|
||||
{ id: "permission_1", sessionID: "ses_other", permission: "file.write", patterns: [], metadata: {}, always: [] } as any,
|
||||
{
|
||||
id: "permission_1",
|
||||
sessionID: "ses_other",
|
||||
permission: "file.write",
|
||||
patterns: [],
|
||||
metadata: {},
|
||||
always: [],
|
||||
} as any,
|
||||
])
|
||||
|
||||
const sender = RemoteSender.create({
|
||||
|
||||
@@ -191,6 +191,29 @@ describe("RemoteWS", () => {
|
||||
expect(server.clients.length).toBe(0)
|
||||
})
|
||||
|
||||
test("onClose callback fires on permanent close", async () => {
|
||||
server = createServer()
|
||||
const codes: number[] = []
|
||||
|
||||
conn = RemoteWS.connect({
|
||||
url: server.url,
|
||||
getToken: async () => "tok",
|
||||
getSessions: () => [],
|
||||
log: nolog(),
|
||||
heartbeat: 60_000,
|
||||
onClose: (code) => codes.push(code),
|
||||
})
|
||||
|
||||
const ws1 = await server.waitForConnect()
|
||||
await settled()
|
||||
|
||||
ws1.close(4401, "unauthorized")
|
||||
await Bun.sleep(100)
|
||||
|
||||
expect(codes).toEqual([4401])
|
||||
expect(conn.connected).toBe(false)
|
||||
})
|
||||
|
||||
test("incoming message delivered to onMessage", async () => {
|
||||
server = createServer()
|
||||
const received: unknown[] = []
|
||||
|
||||
Reference in New Issue
Block a user