From f73c5541fe56196b8b34ab87c045a042ecf95391 Mon Sep 17 00:00:00 2001 From: xingzhou Date: Sat, 11 Jul 2026 22:02:43 +0800 Subject: [PATCH] fix(agents): cancel pending bundle LSP requests (#104110) Co-authored-by: Vincent Koc --- src/agents/agent-bundle-lsp-runtime.test.ts | 68 ++++++++++++- src/agents/agent-bundle-lsp-runtime.ts | 102 ++++++++++++++------ 2 files changed, 140 insertions(+), 30 deletions(-) diff --git a/src/agents/agent-bundle-lsp-runtime.test.ts b/src/agents/agent-bundle-lsp-runtime.test.ts index bc52c4f4cc22..ba3f989d6455 100644 --- a/src/agents/agent-bundle-lsp-runtime.test.ts +++ b/src/agents/agent-bundle-lsp-runtime.test.ts @@ -46,6 +46,7 @@ class MockChildProcess extends EventEmitter { readonly stdout = new PassThrough(); readonly stderr = new PassThrough(); readonly stdin: Writable; + readonly receivedMessages: Record[] = []; constructor( private readonly initializeResponsePrefix = "", @@ -70,13 +71,26 @@ class MockChildProcess extends EventEmitter { private respondToRequest(text: string): void { const body = parseWrittenLspBody(text); - if (!body || typeof body.id !== "number" || typeof body.method !== "string") { + if (!body) { + return; + } + this.receivedMessages.push(body); + if (typeof body.id !== "number" || typeof body.method !== "string") { return; } if (this.respondMethods && !this.respondMethods.has(body.method)) { return; } - const result = body.method === "initialize" ? { capabilities: { hoverProvider: true } } : null; + const result = + body.method === "initialize" + ? { + capabilities: { + hoverProvider: true, + definitionProvider: true, + referencesProvider: true, + }, + } + : null; queueMicrotask(() => { this.stdout.write( `${this.initializeResponsePrefix}${encodeLspMessage({ jsonrpc: "2.0", id: body.id, result })}`, @@ -240,6 +254,56 @@ describe("bundle LSP runtime", () => { await runtime.dispose(); }); + it.each([ + ["lsp_hover_typescript", "textDocument/hover"], + ["lsp_definition_typescript", "textDocument/definition"], + ["lsp_references_typescript", "textDocument/references"], + ])("cancels pending %s requests when the tool signal aborts", async (toolName, method) => { + configureSingleLspServer(); + const child = new MockChildProcess("", new Set(["initialize"])); + spawnMock.mockReturnValue(child); + const { createBundleLspToolRuntime } = await import("./agent-bundle-lsp-runtime.js"); + + const runtime = await createBundleLspToolRuntime({ workspaceDir: "/tmp/workspace" }); + const tool = runtime.tools.find((candidate) => candidate.name === toolName); + if (!tool) { + throw new Error(`expected ${toolName} tool`); + } + const controller = new AbortController(); + const request = tool.execute( + "call-1", + { + uri: "file:///tmp/workspace/index.ts", + line: 0, + character: 0, + }, + controller.signal, + ); + const settled = request.then( + () => "resolved", + () => "rejected", + ); + const lspRequest = child.receivedMessages.find((message) => message.method === method); + + controller.abort(new Error("agent stopped")); + + await expect( + Promise.race([ + settled, + new Promise((resolve) => { + setTimeout(() => resolve("still pending"), 100); + }), + ]), + ).resolves.toBe("rejected"); + expect(child.receivedMessages).toContainEqual({ + jsonrpc: "2.0", + method: "$/cancelRequest", + params: { id: lspRequest?.id }, + }); + + await runtime.dispose(); + }); + it("keeps LSP framing aligned after multibyte messages in the same chunk", async () => { configureSingleLspServer(); const prefix = encodeLspMessage({ diff --git a/src/agents/agent-bundle-lsp-runtime.ts b/src/agents/agent-bundle-lsp-runtime.ts index a28683c42c6f..6424a7168ab3 100644 --- a/src/agents/agent-bundle-lsp-runtime.ts +++ b/src/agents/agent-bundle-lsp-runtime.ts @@ -2,6 +2,7 @@ import { spawn, type ChildProcess } from "node:child_process"; import { normalizeOptionalLowercaseString } from "@openclaw/normalization-core/string-coerce"; import type { OpenClawConfig } from "../config/types.openclaw.js"; +import { createAbortError } from "../infra/abort-signal.js"; import { sanitizeHostExecEnv } from "../infra/host-env-security.js"; import { logDebug, logWarn } from "../logger.js"; import { @@ -39,6 +40,7 @@ type PendingLspRequest = { resolve: (v: unknown) => void; reject: (e: Error) => void; timeout: ReturnType; + dispose: () => void; }; type LspServerCapabilities = { @@ -114,13 +116,22 @@ function rememberLspFailure(session: LspSession, error: Error): void { session.failure ??= error; } +function takePendingLspRequest(session: LspSession, id: number): PendingLspRequest | undefined { + const pending = session.pendingRequests.get(id); + if (!pending) { + return undefined; + } + session.pendingRequests.delete(id); + clearTimeout(pending.timeout); + pending.dispose(); + return pending; +} + function failLspSession(session: LspSession, error: Error): void { rememberLspFailure(session, error); - for (const pending of session.pendingRequests.values()) { - clearTimeout(pending.timeout); - pending.reject(session.failure ?? error); + for (const [id] of session.pendingRequests) { + takePendingLspRequest(session, id)?.reject(session.failure ?? error); } - session.pendingRequests.clear(); } function lspProcessExitError( @@ -205,20 +216,49 @@ function parseLspMessages(buffer: Buffer): { messages: unknown[]; remaining: Buf return { messages, remaining }; } -function sendRequest(session: LspSession, method: string, params?: unknown): Promise { +function lspAbortError(signal?: AbortSignal): Error { + return signal?.reason instanceof Error + ? signal.reason + : createAbortError("LSP request aborted", { cause: signal?.reason }); +} + +function sendRequest( + session: LspSession, + method: string, + params?: unknown, + signal?: AbortSignal, +): Promise { if (session.failure) { return Promise.reject(session.failure); } + if (signal?.aborted) { + return Promise.reject(lspAbortError(signal)); + } const id = ++session.requestId; return new Promise((resolve, reject) => { const timeout = setTimeout(() => { - if (session.pendingRequests.has(id)) { - session.pendingRequests.delete(id); - reject(new Error(`LSP request ${method} timed out`)); - } + takePendingLspRequest(session, id)?.reject(new Error(`LSP request ${method} timed out`)); }, 10_000); timeout.unref?.(); - session.pendingRequests.set(id, { resolve, reject, timeout }); + const onAbort = () => { + const pending = takePendingLspRequest(session, id); + if (!pending) { + return; + } + // Bundle tools share the server process, so cancel only this request. + try { + session.process.stdin?.write( + encodeLspMessage({ jsonrpc: "2.0", method: "$/cancelRequest", params: { id } }), + "utf-8", + ); + } catch { + // Best-effort notification; the local tool promise must still settle. + } + pending.reject(lspAbortError(signal)); + }; + const dispose = () => signal?.removeEventListener("abort", onAbort); + session.pendingRequests.set(id, { resolve, reject, timeout, dispose }); + signal?.addEventListener("abort", onAbort, { once: true }); const message = { jsonrpc: "2.0", id, method, params }; const encoded = encodeLspMessage(message); session.process.stdin?.write(encoded, "utf-8"); @@ -240,10 +280,8 @@ function handleIncomingData(session: LspSession, chunk: Buffer | string) { const record = msg as Record; if ("id" in record && typeof record.id === "number") { - const pending = session.pendingRequests.get(record.id); + const pending = takePendingLspRequest(session, record.id); if (pending) { - session.pendingRequests.delete(record.id); - clearTimeout(pending.timeout); if ("error" in record) { pending.reject(new Error(JSON.stringify(record.error))); } else { @@ -316,11 +354,9 @@ async function disposeSession(session: LspSession) { // best-effort } } - for (const [, pending] of session.pendingRequests) { - clearTimeout(pending.timeout); - pending.reject(new Error("LSP session disposed")); + for (const [id] of session.pendingRequests) { + takePendingLspRequest(session, id)?.reject(new Error("LSP session disposed")); } - session.pendingRequests.clear(); terminateLspProcessTree(session); } @@ -349,12 +385,17 @@ function createLspPositionTool(params: { }, required: ["uri", "line", "character"], }, - execute: async (_toolCallId, input) => { + execute: async (_toolCallId, input, signal) => { const position = input as LspPositionParams; - const result = await sendRequest(params.session, params.method, { - textDocument: { uri: position.uri }, - position: { line: position.line, character: position.character }, - }); + const result = await sendRequest( + params.session, + params.method, + { + textDocument: { uri: position.uri }, + position: { line: position.line, character: position.character }, + }, + signal, + ); return formatLspResult(params.session.serverName, params.resultLabel, result); }, }; @@ -409,18 +450,23 @@ function buildLspTools(session: LspSession): AnyAgentTool[] { }, required: ["uri", "line", "character"], }, - execute: async (_toolCallId, input) => { + execute: async (_toolCallId, input, signal) => { const params = input as { uri: string; line: number; character: number; includeDeclaration?: boolean; }; - const result = await sendRequest(session, "textDocument/references", { - textDocument: { uri: params.uri }, - position: { line: params.line, character: params.character }, - context: { includeDeclaration: params.includeDeclaration ?? true }, - }); + const result = await sendRequest( + session, + "textDocument/references", + { + textDocument: { uri: params.uri }, + position: { line: params.line, character: params.character }, + context: { includeDeclaration: params.includeDeclaration ?? true }, + }, + signal, + ); return formatLspResult(serverLabel, "references", result); }, });