diff --git a/packages/coding-agent/src/dap/client.ts b/packages/coding-agent/src/dap/client.ts index a93e1df9c..ed34af953 100644 --- a/packages/coding-agent/src/dap/client.ts +++ b/packages/coding-agent/src/dap/client.ts @@ -29,32 +29,67 @@ type DapReverseRequestHandler = (args: unknown) => unknown | Promise; const DEFAULT_REQUEST_TIMEOUT_MS = 30_000; -function findHeaderEnd(buffer: Uint8Array): number { - for (let index = 0; index < buffer.length - 3; index += 1) { - if (buffer[index] === 13 && buffer[index + 1] === 10 && buffer[index + 2] === 13 && buffer[index + 3] === 10) { - return index; +// Reused for all full decodes; each decode() resets state, so a single +// instance is safe and avoids per-message TextDecoder allocation. +const MESSAGE_DECODER = new TextDecoder("utf-8"); + +/** + * Locate the `\r\n\r\n` header terminator across the pending chunk list. + * Returns the absolute byte index of the first `\r`, or -1 when not present. + * Equivalent to scanning the contiguous concatenation of the chunks. + */ +function findHeaderEndInChunks(chunks: Buffer[]): number { + let global = 0; + let b0 = -1; + let b1 = -1; + let b2 = -1; + for (const chunk of chunks) { + for (let i = 0; i < chunk.length; i++) { + const b3 = chunk[i]; + if (b0 === 13 && b1 === 10 && b2 === 13 && b3 === 10) { + return global - 3; + } + b0 = b1; + b1 = b2; + b2 = b3; + global++; } } return -1; } -function parseMessage( - buffer: Buffer, -): { message: DapResponseMessage | DapEventMessage | DapRequestMessage; remaining: Buffer } | null { - const headerEndIndex = findHeaderEnd(buffer); - if (headerEndIndex === -1) return null; - const headerText = new TextDecoder().decode(buffer.slice(0, headerEndIndex)); - const contentLengthMatch = headerText.match(/Content-Length: (\d+)/i); - if (!contentLengthMatch) return null; - const contentLength = Number.parseInt(contentLengthMatch[1], 10); - const messageStart = headerEndIndex + 4; - const messageEnd = messageStart + contentLength; - if (buffer.length < messageEnd) return null; - const messageText = new TextDecoder().decode(buffer.subarray(messageStart, messageEnd)); - return { - message: JSON.parse(messageText) as DapResponseMessage | DapEventMessage | DapRequestMessage, - remaining: buffer.subarray(messageEnd), - }; +/** Copy the byte range [from, to) out of the pending chunk list into one Buffer. */ +function copyChunkRange(chunks: Buffer[], from: number, to: number): Buffer { + const out = Buffer.allocUnsafe(to - from); + let global = 0; + let written = 0; + for (const chunk of chunks) { + const chunkEnd = global + chunk.length; + if (chunkEnd > from && global < to) { + const start = Math.max(from, global) - global; + const end = Math.min(to, chunkEnd) - global; + chunk.copy(out, written, start, end); + written += end - start; + } + global = chunkEnd; + if (global >= to) break; + } + return out; +} + +/** Drop the first `count` bytes from the pending chunk list in place. */ +function dropChunkFront(chunks: Buffer[], count: number): void { + let removed = 0; + while (chunks.length > 0) { + const head = chunks[0]; + if (removed + head.length <= count) { + removed += head.length; + chunks.shift(); + } else { + chunks[0] = head.subarray(count - removed); + break; + } + } } async function writeMessage(sink: DapWriteSink, message: DapRequestMessage | DapResponseMessage): Promise { @@ -81,7 +116,7 @@ export class DapClient { readonly #socket?: { end(): void }; #requestSeq = 0; #pendingRequests = new Map(); - #messageBuffer = Buffer.alloc(0); + #messageBuffer: Buffer = Buffer.alloc(0); #isReading = false; #disposed = false; #lastActivity = Date.now(); @@ -416,32 +451,84 @@ export class DapClient { if (this.#isReading) return; this.#isReading = true; const reader = this.#readable.getReader(); + + // Incoming bytes are buffered as a list of chunks and only joined when a + // full message is framed (mirrors the LSP reader) — concatenating the + // accumulator on every read is O(n^2) for messages spanning many reads. + const pendingChunks: Buffer[] = []; + let pendingLen = 0; + if (this.#messageBuffer.length > 0) { + pendingChunks.push(this.#messageBuffer); + pendingLen = this.#messageBuffer.length; + } + try { while (true) { const { done, value } = await reader.read(); if (done) break; - const currentBuffer = Buffer.concat([this.#messageBuffer, value]); - this.#messageBuffer = currentBuffer; - let workingBuffer = currentBuffer; - let parsed = parseMessage(workingBuffer); - while (parsed) { - const { message, remaining } = parsed; - workingBuffer = Buffer.from(remaining); - this.#lastActivity = Date.now(); - if (message.type === "response") { - this.#handleResponse(message); - } else if (message.type === "event") { - await this.#dispatchEvent(message); - } else { - await this.#handleAdapterRequest(message); + + pendingChunks.push(Buffer.from(value)); + pendingLen += value.length; + + // Drain every complete message currently buffered. + while (true) { + const headerEnd = findHeaderEndInChunks(pendingChunks); + if (headerEnd === -1) break; + + const headerText = MESSAGE_DECODER.decode(copyChunkRange(pendingChunks, 0, headerEnd)); + const contentLengthMatch = headerText.match(/Content-Length: (\d+)/i); + if (!contentLengthMatch) { + // Non-protocol bytes (e.g. an adapter printing to stdout). + // Drop past the bogus terminator and resync instead of + // stalling on the same junk header forever. + logger.warn("DAP framing resync: header block without Content-Length", { + adapter: this.adapter.name, + header: headerText.slice(0, 200), + }); + dropChunkFront(pendingChunks, headerEnd + 4); + pendingLen -= headerEnd + 4; + continue; + } + + const contentLength = Number.parseInt(contentLengthMatch[1], 10); + const messageStart = headerEnd + 4; // Skip \r\n\r\n + const messageEnd = messageStart + contentLength; + if (pendingLen < messageEnd) break; + + const messageText = MESSAGE_DECODER.decode(copyChunkRange(pendingChunks, messageStart, messageEnd)); + dropChunkFront(pendingChunks, messageEnd); + pendingLen -= messageEnd; + this.#lastActivity = Date.now(); + + // A malformed message must not kill the reader — later + // messages are still well-framed. + try { + const message = JSON.parse(messageText) as DapResponseMessage | DapEventMessage | DapRequestMessage; + if (message.type === "response") { + this.#handleResponse(message); + } else if (message.type === "event") { + await this.#dispatchEvent(message); + } else { + await this.#handleAdapterRequest(message); + } + } catch (error) { + logger.warn("DAP message handling failed", { + adapter: this.adapter.name, + error: toErrorMessage(error), + }); } - parsed = parseMessage(workingBuffer); } - this.#messageBuffer = workingBuffer; } } catch (error) { this.#rejectPendingRequests(new Error(`DAP connection closed: ${toErrorMessage(error)}`)); } finally { + // Persist any unparsed remainder so a restarted reader resumes mid-message. + this.#messageBuffer = + pendingChunks.length === 0 + ? Buffer.alloc(0) + : pendingChunks.length === 1 + ? pendingChunks[0] + : Buffer.concat(pendingChunks, pendingLen); reader.releaseLock(); this.#isReading = false; } diff --git a/packages/coding-agent/src/dap/session.ts b/packages/coding-agent/src/dap/session.ts index 57be0afc1..82f9ba3e2 100644 --- a/packages/coding-agent/src/dap/session.ts +++ b/packages/coding-agent/src/dap/session.ts @@ -76,8 +76,14 @@ interface DapSession { functionBreakpoints: DapFunctionBreakpointRecord[]; instructionBreakpoints: DapInstructionBreakpoint[]; dataBreakpoints: DapDataBreakpoint[]; - output: string; + /** Serializes breakpoint mutations — see #serializeBreakpointMutation. */ + breakpointMutationQueue: Promise; + /** Recent output chunks; trimmed from the front when over MAX_OUTPUT_BYTES. */ + outputChunks: string[]; + /** Cumulative bytes of output ever received (reported in summaries). */ outputBytes: number; + /** Bytes currently buffered in outputChunks. */ + outputBufferedBytes: number; outputTruncated: boolean; stop: DapStopLocation; threads: DapThread[]; @@ -175,10 +181,31 @@ function normalizePath(filePath: string): string { function truncateOutput(session: DapSession, output: string): void { if (!output) return; - session.output += output; - session.outputBytes += Buffer.byteLength(output, "utf-8"); - while (Buffer.byteLength(session.output, "utf-8") > MAX_OUTPUT_BYTES) { - session.output = session.output.slice(Math.min(1024, session.output.length)); + const bytes = Buffer.byteLength(output, "utf-8"); + session.outputChunks.push(output); + session.outputBytes += bytes; + session.outputBufferedBytes += bytes; + // Trim whole chunks from the front, but only while the remainder still + // holds a full MAX_OUTPUT_BYTES tail — dropping the front chunk whenever + // the total exceeded the cap could retain far less than the cap (e.g. + // [120KB, 10KB] would keep only 10KB). Recomputing one big string's byte + // length per 1KB trim iteration was O(n^2) inside the event dispatch loop. + while (session.outputChunks.length > 1) { + const frontBytes = Buffer.byteLength(session.outputChunks[0], "utf-8"); + if (session.outputBufferedBytes - frontBytes < MAX_OUTPUT_BYTES) break; + session.outputChunks.shift(); + session.outputBufferedBytes -= frontBytes; + session.outputTruncated = true; + } + if (session.outputBufferedBytes > MAX_OUTPUT_BYTES) { + // Byte-slice the front chunk's head so exactly the cap remains (a torn + // code point at the cut decodes as U+FFFD, acceptable for log output). + const front = session.outputChunks[0]; + const frontBytes = Buffer.byteLength(front, "utf-8"); + const excess = session.outputBufferedBytes - MAX_OUTPUT_BYTES; + const kept = Buffer.from(front, "utf-8").subarray(excess).toString("utf-8"); + session.outputChunks[0] = kept; + session.outputBufferedBytes += Buffer.byteLength(kept, "utf-8") - frontBytes; session.outputTruncated = true; } } @@ -368,6 +395,26 @@ export class DapSessionManager { } } + /** + * Serialize breakpoint mutations per session: every mutator does a + * read-modify-write of session state around an await, and the adapter-side + * set*Breakpoints request replaces the whole list — concurrent mutations + * would silently drop each other's breakpoints on both sides. + */ + #serializeBreakpointMutation(session: DapSession, mutate: () => Promise, signal?: AbortSignal): Promise { + const run = session.breakpointMutationQueue.then(() => { + // A mutation can sit behind several queued 30s predecessors; honor a + // caller abort at dequeue instead of running a request nobody awaits. + if (signal?.aborted) throw signal.reason instanceof Error ? signal.reason : new Error("Aborted"); + return mutate(); + }); + session.breakpointMutationQueue = run.then( + () => undefined, + () => undefined, + ); + return run; + } + async setBreakpoint( file: string, line: number, @@ -376,99 +423,123 @@ export class DapSessionManager { timeoutMs: number = 30_000, ) { const session = this.#touchActiveSession(); - const sourcePath = normalizePath(file); - const current = [...(session.breakpoints.get(sourcePath) ?? [])]; - const deduped = current.filter(entry => entry.line !== line); - deduped.push({ verified: false, line, condition }); - deduped.sort((left, right) => left.line - right.line); - const response = await this.#sendRequestWithConfig<{ breakpoints?: DapBreakpoint[] }>( + return this.#serializeBreakpointMutation( session, - "setBreakpoints", - { - source: { path: sourcePath, name: path.basename(sourcePath) }, - breakpoints: deduped.map(entry => ({ - line: entry.line, - ...(entry.condition ? { condition: entry.condition } : {}), - })), + async () => { + const sourcePath = normalizePath(file); + const current = [...(session.breakpoints.get(sourcePath) ?? [])]; + const deduped = current.filter(entry => entry.line !== line); + deduped.push({ verified: false, line, condition }); + deduped.sort((left, right) => left.line - right.line); + const response = await this.#sendRequestWithConfig<{ breakpoints?: DapBreakpoint[] }>( + session, + "setBreakpoints", + { + source: { path: sourcePath, name: path.basename(sourcePath) }, + breakpoints: deduped.map(entry => ({ + line: entry.line, + ...(entry.condition ? { condition: entry.condition } : {}), + })), + }, + signal, + timeoutMs, + ); + session.breakpoints.set(sourcePath, this.#mapSourceBreakpoints(deduped, response?.breakpoints)); + return { + snapshot: buildSummary(session), + breakpoints: session.breakpoints.get(sourcePath) ?? [], + sourcePath, + }; }, signal, - timeoutMs, ); - session.breakpoints.set(sourcePath, this.#mapSourceBreakpoints(deduped, response?.breakpoints)); - return { - snapshot: buildSummary(session), - breakpoints: session.breakpoints.get(sourcePath) ?? [], - sourcePath, - }; } async removeBreakpoint(file: string, line: number, signal?: AbortSignal, timeoutMs: number = 30_000) { const session = this.#touchActiveSession(); - const sourcePath = normalizePath(file); - const current = [...(session.breakpoints.get(sourcePath) ?? [])].filter(entry => entry.line !== line); - const response = await this.#sendRequestWithConfig<{ breakpoints?: DapBreakpoint[] }>( + return this.#serializeBreakpointMutation( session, - "setBreakpoints", - { - source: { path: sourcePath, name: path.basename(sourcePath) }, - breakpoints: current.map(entry => ({ - line: entry.line, - ...(entry.condition ? { condition: entry.condition } : {}), - })), + async () => { + const sourcePath = normalizePath(file); + const current = [...(session.breakpoints.get(sourcePath) ?? [])].filter(entry => entry.line !== line); + const response = await this.#sendRequestWithConfig<{ breakpoints?: DapBreakpoint[] }>( + session, + "setBreakpoints", + { + source: { path: sourcePath, name: path.basename(sourcePath) }, + breakpoints: current.map(entry => ({ + line: entry.line, + ...(entry.condition ? { condition: entry.condition } : {}), + })), + }, + signal, + timeoutMs, + ); + if (current.length === 0) { + session.breakpoints.delete(sourcePath); + } else { + session.breakpoints.set(sourcePath, this.#mapSourceBreakpoints(current, response?.breakpoints)); + } + return { + snapshot: buildSummary(session), + breakpoints: session.breakpoints.get(sourcePath) ?? [], + sourcePath, + }; }, signal, - timeoutMs, ); - if (current.length === 0) { - session.breakpoints.delete(sourcePath); - } else { - session.breakpoints.set(sourcePath, this.#mapSourceBreakpoints(current, response?.breakpoints)); - } - return { - snapshot: buildSummary(session), - breakpoints: session.breakpoints.get(sourcePath) ?? [], - sourcePath, - }; } async setFunctionBreakpoint(name: string, condition?: string, signal?: AbortSignal, timeoutMs: number = 30_000) { const session = this.#touchActiveSession(); - const current = session.functionBreakpoints.filter(entry => entry.name !== name); - current.push({ verified: false, name, condition }); - current.sort((left, right) => left.name.localeCompare(right.name)); - const response = await this.#sendRequestWithConfig<{ breakpoints?: DapBreakpoint[] }>( + return this.#serializeBreakpointMutation( session, - "setFunctionBreakpoints", - { - breakpoints: current.map(entry => ({ - name: entry.name, - ...(entry.condition ? { condition: entry.condition } : {}), - })), + async () => { + const current = session.functionBreakpoints.filter(entry => entry.name !== name); + current.push({ verified: false, name, condition }); + current.sort((left, right) => left.name.localeCompare(right.name)); + const response = await this.#sendRequestWithConfig<{ breakpoints?: DapBreakpoint[] }>( + session, + "setFunctionBreakpoints", + { + breakpoints: current.map(entry => ({ + name: entry.name, + ...(entry.condition ? { condition: entry.condition } : {}), + })), + }, + signal, + timeoutMs, + ); + session.functionBreakpoints = this.#mapFunctionBreakpoints(current, response?.breakpoints); + return { snapshot: buildSummary(session), breakpoints: session.functionBreakpoints }; }, signal, - timeoutMs, ); - session.functionBreakpoints = this.#mapFunctionBreakpoints(current, response?.breakpoints); - return { snapshot: buildSummary(session), breakpoints: session.functionBreakpoints }; } async removeFunctionBreakpoint(name: string, signal?: AbortSignal, timeoutMs: number = 30_000) { const session = this.#touchActiveSession(); - const current = session.functionBreakpoints.filter(entry => entry.name !== name); - const response = await this.#sendRequestWithConfig<{ breakpoints?: DapBreakpoint[] }>( + return this.#serializeBreakpointMutation( session, - "setFunctionBreakpoints", - { - breakpoints: current.map(entry => ({ - name: entry.name, - ...(entry.condition ? { condition: entry.condition } : {}), - })), + async () => { + const current = session.functionBreakpoints.filter(entry => entry.name !== name); + const response = await this.#sendRequestWithConfig<{ breakpoints?: DapBreakpoint[] }>( + session, + "setFunctionBreakpoints", + { + breakpoints: current.map(entry => ({ + name: entry.name, + ...(entry.condition ? { condition: entry.condition } : {}), + })), + }, + signal, + timeoutMs, + ); + session.functionBreakpoints = this.#mapFunctionBreakpoints(current, response?.breakpoints); + return { snapshot: buildSummary(session), breakpoints: session.functionBreakpoints }; }, signal, - timeoutMs, ); - session.functionBreakpoints = this.#mapFunctionBreakpoints(current, response?.breakpoints); - return { snapshot: buildSummary(session), breakpoints: session.functionBreakpoints }; } async setInstructionBreakpoint( @@ -480,31 +551,37 @@ export class DapSessionManager { timeoutMs: number = 30_000, ) { const session = this.#touchActiveSession(); - const current = session.instructionBreakpoints.filter( - entry => entry.instructionReference !== instructionReference || entry.offset !== offset, - ); - current.push({ instructionReference, offset, condition, hitCondition }); - current.sort((left, right) => { - const referenceOrder = left.instructionReference.localeCompare(right.instructionReference); - if (referenceOrder !== 0) { - return referenceOrder; - } - return (left.offset ?? 0) - (right.offset ?? 0); - }); - const response = await this.#sendRequestWithConfig<{ breakpoints?: DapBreakpoint[] }>( + return this.#serializeBreakpointMutation( session, - "setInstructionBreakpoints", - { - breakpoints: current, - } satisfies DapSetInstructionBreakpointsArguments, + async () => { + const current = session.instructionBreakpoints.filter( + entry => entry.instructionReference !== instructionReference || entry.offset !== offset, + ); + current.push({ instructionReference, offset, condition, hitCondition }); + current.sort((left, right) => { + const referenceOrder = left.instructionReference.localeCompare(right.instructionReference); + if (referenceOrder !== 0) { + return referenceOrder; + } + return (left.offset ?? 0) - (right.offset ?? 0); + }); + const response = await this.#sendRequestWithConfig<{ breakpoints?: DapBreakpoint[] }>( + session, + "setInstructionBreakpoints", + { + breakpoints: current, + } satisfies DapSetInstructionBreakpointsArguments, + signal, + timeoutMs, + ); + session.instructionBreakpoints = current; + return { + snapshot: buildSummary(session), + breakpoints: this.#mapInstructionBreakpoints(current, response?.breakpoints), + }; + }, signal, - timeoutMs, ); - session.instructionBreakpoints = current; - return { - snapshot: buildSummary(session), - breakpoints: this.#mapInstructionBreakpoints(current, response?.breakpoints), - }; } async removeInstructionBreakpoint( @@ -514,29 +591,35 @@ export class DapSessionManager { timeoutMs: number = 30_000, ) { const session = this.#touchActiveSession(); - const current = session.instructionBreakpoints.filter(entry => { - if (entry.instructionReference !== instructionReference) { - return true; - } - if (offset === undefined) { - return false; - } - return entry.offset !== offset; - }); - const response = await this.#sendRequestWithConfig<{ breakpoints?: DapBreakpoint[] }>( + return this.#serializeBreakpointMutation( session, - "setInstructionBreakpoints", - { - breakpoints: current, - } satisfies DapSetInstructionBreakpointsArguments, + async () => { + const current = session.instructionBreakpoints.filter(entry => { + if (entry.instructionReference !== instructionReference) { + return true; + } + if (offset === undefined) { + return false; + } + return entry.offset !== offset; + }); + const response = await this.#sendRequestWithConfig<{ breakpoints?: DapBreakpoint[] }>( + session, + "setInstructionBreakpoints", + { + breakpoints: current, + } satisfies DapSetInstructionBreakpointsArguments, + signal, + timeoutMs, + ); + session.instructionBreakpoints = current; + return { + snapshot: buildSummary(session), + breakpoints: this.#mapInstructionBreakpoints(current, response?.breakpoints), + }; + }, signal, - timeoutMs, ); - session.instructionBreakpoints = current; - return { - snapshot: buildSummary(session), - breakpoints: this.#mapInstructionBreakpoints(current, response?.breakpoints), - }; } async dataBreakpointInfo( @@ -570,42 +653,54 @@ export class DapSessionManager { timeoutMs: number = 30_000, ) { const session = this.#touchActiveSession(); - const current = session.dataBreakpoints.filter(entry => entry.dataId !== dataId); - current.push({ dataId, accessType, condition, hitCondition }); - current.sort((left, right) => left.dataId.localeCompare(right.dataId)); - const response = await this.#sendRequestWithConfig<{ breakpoints?: DapBreakpoint[] }>( + return this.#serializeBreakpointMutation( session, - "setDataBreakpoints", - { - breakpoints: current, - } satisfies DapSetDataBreakpointsArguments, + async () => { + const current = session.dataBreakpoints.filter(entry => entry.dataId !== dataId); + current.push({ dataId, accessType, condition, hitCondition }); + current.sort((left, right) => left.dataId.localeCompare(right.dataId)); + const response = await this.#sendRequestWithConfig<{ breakpoints?: DapBreakpoint[] }>( + session, + "setDataBreakpoints", + { + breakpoints: current, + } satisfies DapSetDataBreakpointsArguments, + signal, + timeoutMs, + ); + session.dataBreakpoints = current; + return { + snapshot: buildSummary(session), + breakpoints: this.#mapDataBreakpoints(current, response?.breakpoints), + }; + }, signal, - timeoutMs, ); - session.dataBreakpoints = current; - return { - snapshot: buildSummary(session), - breakpoints: this.#mapDataBreakpoints(current, response?.breakpoints), - }; } async removeDataBreakpoint(dataId: string, signal?: AbortSignal, timeoutMs: number = 30_000) { const session = this.#touchActiveSession(); - const current = session.dataBreakpoints.filter(entry => entry.dataId !== dataId); - const response = await this.#sendRequestWithConfig<{ breakpoints?: DapBreakpoint[] }>( + return this.#serializeBreakpointMutation( session, - "setDataBreakpoints", - { - breakpoints: current, - } satisfies DapSetDataBreakpointsArguments, + async () => { + const current = session.dataBreakpoints.filter(entry => entry.dataId !== dataId); + const response = await this.#sendRequestWithConfig<{ breakpoints?: DapBreakpoint[] }>( + session, + "setDataBreakpoints", + { + breakpoints: current, + } satisfies DapSetDataBreakpointsArguments, + signal, + timeoutMs, + ); + session.dataBreakpoints = current; + return { + snapshot: buildSummary(session), + breakpoints: this.#mapDataBreakpoints(current, response?.breakpoints), + }; + }, signal, - timeoutMs, ); - session.dataBreakpoints = current; - return { - snapshot: buildSummary(session), - breakpoints: this.#mapDataBreakpoints(current, response?.breakpoints), - }; } async disassemble( @@ -756,21 +851,25 @@ export class DapSessionManager { async pause(signal?: AbortSignal, timeoutMs: number = 30_000): Promise { const session = this.#touchActiveSession(); - if (session.status === "stopped") { + // status is mutated by the event reader between awaits; check through a + // closure so TS does not carry stale narrowing from the early return. + const isStopped = () => session.status === "stopped"; + if (isStopped()) { return buildSummary(session); } const threadId = await this.#resolveThreadId(session, signal, timeoutMs); + // Subscribe BEFORE sending pause: the stopped event can arrive in the + // same chunk as the response and would otherwise be dispatched before + // the waiter subscribes, burning the whole timeout. + const stoppedPromise = session.client.waitForEvent("stopped", undefined, signal, timeoutMs); + stoppedPromise.catch(() => {}); await this.#sendRequestWithConfig(session, "pause", { threadId } satisfies DapPauseArguments, signal, timeoutMs); - // The stopped event may already have been processed by #handleStoppedEvent - // between the request and here. Wait for it, but tolerate timeout if the - // session already transitioned. - try { - await untilAborted( - signal, - session.client.waitForEvent("stopped", undefined, signal, timeoutMs), - ); - } catch { - // Timeout or abort — report current state regardless + if (!isStopped()) { + try { + await untilAborted(signal, stoppedPromise); + } catch { + // Timeout or abort — report current state regardless + } } return buildSummary(session); } @@ -884,16 +983,16 @@ export class DapSessionManager { getOutput(limitBytes?: number): DapOutputSnapshot { const session = this.#touchActiveSession(); - if (!limitBytes || limitBytes <= 0 || Buffer.byteLength(session.output, "utf-8") <= limitBytes) { - return { snapshot: buildSummary(session), output: session.output }; + const output = session.outputChunks.join(""); + if (!limitBytes || limitBytes <= 0 || session.outputBufferedBytes <= limitBytes) { + return { snapshot: buildSummary(session), output }; } - let sliceStart = session.output.length; - let remaining = limitBytes; - while (sliceStart > 0 && remaining > 0) { - sliceStart -= 1; - remaining -= Buffer.byteLength(session.output[sliceStart] ?? "", "utf-8"); + // Byte-slice the tail once; a torn code point at the cut decodes as U+FFFD. + const buffer = Buffer.from(output, "utf-8"); + if (buffer.length <= limitBytes) { + return { snapshot: buildSummary(session), output }; } - return { snapshot: buildSummary(session), output: session.output.slice(sliceStart) }; + return { snapshot: buildSummary(session), output: buffer.subarray(buffer.length - limitBytes).toString("utf-8") }; } async terminate(signal?: AbortSignal, timeoutMs: number = 30_000): Promise { @@ -973,8 +1072,10 @@ export class DapSessionManager { functionBreakpoints: [], instructionBreakpoints: [], dataBreakpoints: [], - output: "", + breakpointMutationQueue: Promise.resolve(), + outputChunks: [], outputBytes: 0, + outputBufferedBytes: 0, outputTruncated: false, stop: {}, threads: [], diff --git a/packages/coding-agent/src/lsp/client.ts b/packages/coding-agent/src/lsp/client.ts index b0e5c5069..a2b6293e1 100644 --- a/packages/coding-agent/src/lsp/client.ts +++ b/packages/coding-agent/src/lsp/client.ts @@ -22,6 +22,10 @@ const clients = new Map(); const clientLocks = new Map>(); const fileOperationLocks = new Map>(); +/** Negative cache of recent init failures so a broken server fails fast instead of re-spawning per call. */ +const INIT_FAILURE_BACKOFF_MS = 3 * 60 * 1000; +const initFailures = new Map(); + // Idle timeout configuration (disabled by default) let idleTimeoutMs: number | null = null; let idleCheckInterval: NodeJS.Timeout | null = null; @@ -295,7 +299,18 @@ async function startMessageReader(client: LspClient): Promise { const headerText = MESSAGE_DECODER.decode(copyChunkRange(pendingChunks, 0, headerEnd)); const contentLengthMatch = headerText.match(/Content-Length: (\d+)/i); - if (!contentLengthMatch) break; + if (!contentLengthMatch) { + // Non-protocol bytes on stdout (e.g. a wrapper script printing). + // Drop past the bogus terminator and resync instead of stalling + // on the same junk header forever. + logger.warn("LSP framing resync: header block without Content-Length", { + server: client.name, + header: headerText.slice(0, 200), + }); + dropChunkFront(pendingChunks, headerEnd + 4); + pendingLen -= headerEnd + 4; + continue; + } const contentLength = Number.parseInt(contentLengthMatch[1], 10); const messageStart = headerEnd + 4; // Skip \r\n\r\n @@ -303,44 +318,54 @@ async function startMessageReader(client: LspClient): Promise { if (pendingLen < messageEnd) break; const messageText = MESSAGE_DECODER.decode(copyChunkRange(pendingChunks, messageStart, messageEnd)); - const message: LspJsonRpcResponse | LspJsonRpcNotification = JSON.parse(messageText); dropChunkFront(pendingChunks, messageEnd); pendingLen -= messageEnd; - // Route message - if ("id" in message && message.id !== undefined) { - // Response to a request - const pending = client.pendingRequests.get(message.id); - if (pending) { - client.pendingRequests.delete(message.id); - if ("error" in message && message.error) { - pending.reject(new Error(`LSP error: ${message.error.message}`)); - } else { - pending.resolve(message.result); + // A malformed message or a throwing server-request handler must not + // kill the reader — later messages are still well-framed. + try { + const message: LspJsonRpcResponse | LspJsonRpcNotification = JSON.parse(messageText); + + // Route message + if ("id" in message && message.id !== undefined) { + // Response to a request + const pending = client.pendingRequests.get(message.id); + if (pending) { + client.pendingRequests.delete(message.id); + if ("error" in message && message.error) { + pending.reject(new Error(`LSP error: ${message.error.message}`)); + } else { + pending.resolve(message.result); + } + } else if ("method" in message) { + await handleServerRequest(client, message as LspJsonRpcRequest); } } else if ("method" in message) { - await handleServerRequest(client, message as LspJsonRpcRequest); - } - } else if ("method" in message) { - // Server notification - if (message.method === "textDocument/publishDiagnostics" && message.params) { - const params = message.params as PublishDiagnosticsParams; - client.diagnostics.set(params.uri, { - diagnostics: params.diagnostics, - version: params.version ?? null, - }); - client.diagnosticsVersion += 1; - } else if (message.method === "$/progress" && message.params) { - const params = message.params as { token: string | number; value?: { kind?: string } }; - if (params.value?.kind === "begin") { - client.activeProgressTokens.add(params.token); - } else if (params.value?.kind === "end") { - client.activeProgressTokens.delete(params.token); - if (client.activeProgressTokens.size === 0) { - client.resolveProjectLoaded(); + // Server notification + if (message.method === "textDocument/publishDiagnostics" && message.params) { + const params = message.params as PublishDiagnosticsParams; + client.diagnostics.set(params.uri, { + diagnostics: params.diagnostics, + version: params.version ?? null, + }); + client.diagnosticsVersion += 1; + } else if (message.method === "$/progress" && message.params) { + const params = message.params as { token: string | number; value?: { kind?: string } }; + if (params.value?.kind === "begin") { + client.activeProgressTokens.add(params.token); + } else if (params.value?.kind === "end") { + client.activeProgressTokens.delete(params.token); + if (client.activeProgressTokens.size === 0) { + client.resolveProjectLoaded(); + } } } } + } catch (err) { + logger.warn("LSP message handling failed", { + server: client.name, + error: err instanceof Error ? err.message : String(err), + }); } } } @@ -360,6 +385,22 @@ async function startMessageReader(client: LspClient): Promise { : Buffer.concat(pendingChunks, pendingLen); reader.releaseLock(); client.isReading = false; + // Reader exited while the server process is still alive (unrecoverable + // read error or bad stream state): nothing will route responses anymore, + // so tear the client down — the next call respawns instead of timing out. + if (client.proc.exitCode === null) { + client.status = "error"; + if (clients.get(client.name) === client) { + clients.delete(client.name); + } + const teardownErr = new Error("LSP reader stopped; client torn down"); + for (const pending of client.pendingRequests.values()) { + pending.reject(teardownErr); + } + client.pendingRequests.clear(); + client.resolveProjectLoaded(); + client.proc.kill(); + } } } @@ -565,6 +606,16 @@ export async function getOrCreateClient(config: ServerConfig, cwd: string, initT return existingLock; } + // Fail fast on a recent deterministic init failure instead of re-spawning + // a broken server (and paying its full init wait) on every call. + const recentFailure = initFailures.get(key); + if (recentFailure) { + if (Date.now() - recentFailure.at < INIT_FAILURE_BACKOFF_MS) { + throw new Error(`LSP server ${config.command} failed to initialize recently: ${recentFailure.message}`); + } + initFailures.delete(key); + } + // Create new client with lock const clientPromise = (async () => { const baseCommand = config.resolvedCommand ?? config.command; @@ -605,18 +656,18 @@ export async function getOrCreateClient(config: ServerConfig, cwd: string, initT pendingRequests: new Map(), messageBuffer: new Uint8Array(0), isReading: false, + status: "connecting", lastActivity: Date.now(), writeQueue: Promise.resolve(), activeProgressTokens: new Set(), projectLoaded, resolveProjectLoaded, }; - clients.set(key, client); // Register crash recovery - remove client on process exit proc.exited.then(() => { - clients.delete(key); - clientLocks.delete(key); + if (clients.get(key) === client) clients.delete(key); + if (clientLocks.get(key) === clientPromise) clientLocks.delete(key); client.resolveProjectLoaded(); // Reject any pending requests — the server is gone, they will never complete. @@ -669,12 +720,26 @@ export async function getOrCreateClient(config: ServerConfig, cwd: string, initT // Send initialized notification await sendNotification(client, "initialized", {}); + client.status = "ready"; + // Publish only after init succeeds: pre-init clients are reachable + // solely through clientLocks, so concurrent callers (warmup vs first + // tool call) wait for init instead of using an unacknowledged client. + clients.set(key, client); + initFailures.delete(key); return client; } catch (err) { // Clean up on initialization failure - clients.delete(key); - clientLocks.delete(key); + client.status = "error"; + if (clients.get(key) === client) clients.delete(key); proc.kill(); + const message = err instanceof Error ? err.message : String(err); + // Negative-cache deterministic failures. Timeouts under a + // caller-shortened deadline (warmup/writethrough) are not cached — + // the server may simply be slow and a later call with the full + // deadline can still succeed. + if (!(initTimeoutMs !== undefined && message.includes("timed out"))) { + initFailures.set(key, { at: Date.now(), message }); + } throw err; } finally { clientLocks.delete(key); @@ -1067,7 +1132,22 @@ export async function sendNotification(client: LspClient, method: string, params export async function shutdownAll(): Promise { const clientsToShutdown = Array.from(clients.values()); clients.clear(); - await Promise.allSettled(clientsToShutdown.map(client => shutdownClientInstance(client))); + // Mid-initialize clients live only in clientLocks (publication is deferred + // until init succeeds) — without this, their server processes outlive + // shutdown. Failed init promises already cleaned up after themselves. + const pendingClients = Array.from(clientLocks.values()); + clientLocks.clear(); + const seen = new Set(clientsToShutdown); + await Promise.allSettled([ + ...clientsToShutdown.map(client => shutdownClientInstance(client)), + ...pendingClients.map(pending => + pending.then(client => { + if (seen.has(client)) return; + seen.add(client); + return shutdownClientInstance(client); + }), + ), + ]); } /** Status of an LSP server */ @@ -1084,7 +1164,7 @@ export interface LspServerStatus { export function getActiveClients(): LspServerStatus[] { return Array.from(clients.values()).map(client => ({ name: client.config.command, - status: "ready" as const, + status: client.status, fileTypes: client.config.fileTypes, })); } diff --git a/packages/coding-agent/src/lsp/clients/biome-client.ts b/packages/coding-agent/src/lsp/clients/biome-client.ts index 2ebb99de6..82bd497a6 100644 --- a/packages/coding-agent/src/lsp/clients/biome-client.ts +++ b/packages/coding-agent/src/lsp/clients/biome-client.ts @@ -3,6 +3,7 @@ * Uses Biome's CLI with JSON output instead of LSP (which has stale diagnostics issues). */ import path from "node:path"; +import { logger } from "@oh-my-pi/pi-utils"; import type { Diagnostic, DiagnosticSeverity, LinterClient, ServerConfig } from "../../lsp/types"; // ============================================================================= @@ -29,17 +30,23 @@ interface BiomeDiagnostic { // ============================================================================= /** - * Convert byte offset to line:column using source code. + * Convert byte offsets to line:column positions in a single pass over the source. */ -function offsetToPosition(source: string, offset: number): { line: number; column: number } { +function offsetsToPositions(source: string, offsets: number[]): Map { + const sorted = [...new Set(offsets)].sort((a, b) => a - b); + const result = new Map(); let line = 1; let column = 1; let byteIndex = 0; + let next = 0; for (const ch of source) { - const byteLen = Buffer.byteLength(ch); - if (byteIndex + byteLen > offset) { - break; + if (next >= sorted.length) break; + const cp = ch.codePointAt(0) as number; + const byteLen = cp < 0x80 ? 1 : cp < 0x800 ? 2 : cp < 0x10000 ? 3 : 4; + while (next < sorted.length && byteIndex + byteLen > sorted[next]) { + result.set(sorted[next], { line, column }); + next++; } if (ch === "\n") { line++; @@ -50,7 +57,13 @@ function offsetToPosition(source: string, offset: number): { line: number; colum byteIndex += byteLen; } - return { line, column }; + // Offsets at or past end-of-file map to the final position. + while (next < sorted.length) { + result.set(sorted[next], { line, column }); + next++; + } + + return result; } /** @@ -98,6 +111,16 @@ async function runBiome( } } +// Surface broken-binary / CLI failures once instead of silently reporting +// "no diagnostics" forever (and instead of spamming every writethrough). +const reportedBiomeFailures = new Set(); + +function warnBiomeOnce(key: string, message: string, meta: Record): void { + if (reportedBiomeFailures.has(key)) return; + reportedBiomeFailures.add(key); + logger.warn(message, meta); +} + // ============================================================================= // Biome Client // ============================================================================= @@ -137,6 +160,16 @@ export class BiomeClient implements LinterClient { // Run biome lint with JSON reporter const result = await runBiome(["lint", "--reporter=json", filePath], this.cwd, this.config.resolvedCommand); + // Biome exits non-zero when diagnostics are found, so only an empty + // stdout signals an actual run failure (missing binary, CLI error). + if (!result.success && result.stdout.trim().length === 0) { + warnBiomeOnce(`run:${this.cwd}`, "Biome lint failed; reporting no diagnostics", { + cwd: this.cwd, + stderr: result.stderr.slice(0, 500), + }); + return []; + } + return this.#parseJsonOutput(result.stdout, filePath); } @@ -146,51 +179,80 @@ export class BiomeClient implements LinterClient { #parseJsonOutput(jsonOutput: string, targetFile: string): Diagnostic[] { const diagnostics: Diagnostic[] = []; + let parsed: BiomeJsonOutput; try { - const parsed: BiomeJsonOutput = JSON.parse(jsonOutput); + parsed = JSON.parse(jsonOutput); + } catch { + warnBiomeOnce(`parse:${this.cwd}`, "Failed to parse Biome JSON output; reporting no diagnostics", { + cwd: this.cwd, + file: targetFile, + }); + return diagnostics; + } - for (const diag of parsed.diagnostics) { - const location = diag.location; - if (!location?.path?.file) continue; + const target = path.resolve(targetFile); + const relevant: BiomeDiagnostic[] = []; + // Batch all span offsets per source text so each source is scanned once + // instead of twice per diagnostic. + const offsetsBySource = new Map(); + for (const diag of parsed.diagnostics ?? []) { + const location = diag.location; + if (!location?.path?.file) continue; - // Resolve file path - const diagFile = path.isAbsolute(location.path.file) - ? location.path.file - : path.join(this.cwd, location.path.file); + // Resolve file path + const diagFile = path.isAbsolute(location.path.file) + ? location.path.file + : path.join(this.cwd, location.path.file); - // Only include diagnostics for the target file - if (path.resolve(diagFile) !== path.resolve(targetFile)) { - continue; - } + // Only include diagnostics for the target file + if (path.resolve(diagFile) !== target) { + continue; + } - // Convert byte offset to line:column - let startLine = 1; - let startColumn = 1; - let endLine = 1; - let endColumn = 1; + relevant.push(diag); + if (location.span && location.sourceCode) { + const offsets = offsetsBySource.get(location.sourceCode); + if (offsets) offsets.push(location.span[0], location.span[1]); + else offsetsBySource.set(location.sourceCode, [location.span[0], location.span[1]]); + } + } - if (location.span && location.sourceCode) { - const startPos = offsetToPosition(location.sourceCode, location.span[0]); - const endPos = offsetToPosition(location.sourceCode, location.span[1]); + const positionsBySource = new Map>(); + for (const [source, offsets] of offsetsBySource) { + positionsBySource.set(source, offsetsToPositions(source, offsets)); + } + + for (const diag of relevant) { + const location = diag.location; + let startLine = 1; + let startColumn = 1; + let endLine = 1; + let endColumn = 1; + + if (location?.span && location.sourceCode) { + const positions = positionsBySource.get(location.sourceCode); + const startPos = positions?.get(location.span[0]); + const endPos = positions?.get(location.span[1]); + if (startPos) { startLine = startPos.line; startColumn = startPos.column; + } + if (endPos) { endLine = endPos.line; endColumn = endPos.column; } - - diagnostics.push({ - range: { - start: { line: startLine - 1, character: startColumn - 1 }, - end: { line: endLine - 1, character: endColumn - 1 }, - }, - severity: parseSeverity(diag.severity), - message: diag.description, - source: "biome", - code: diag.category, - }); } - } catch { - // JSON parse failed, return empty + + diagnostics.push({ + range: { + start: { line: startLine - 1, character: startColumn - 1 }, + end: { line: endLine - 1, character: endColumn - 1 }, + }, + severity: parseSeverity(diag.severity), + message: diag.description, + source: "biome", + code: diag.category, + }); } return diagnostics; diff --git a/packages/coding-agent/src/lsp/edits.ts b/packages/coding-agent/src/lsp/edits.ts index 78c96b262..d038d84bf 100644 --- a/packages/coding-agent/src/lsp/edits.ts +++ b/packages/coding-agent/src/lsp/edits.ts @@ -24,27 +24,7 @@ import { uriToFile } from "./utils"; */ export function applyTextEditsToString(content: string, edits: TextEdit[]): string { const lines = content.split("\n"); - - // Sort edits in reverse order (bottom-to-top, right-to-left) - const sortedEdits = [...edits].sort((a, b) => { - if (a.range.start.line !== b.range.start.line) { - return b.range.start.line - a.range.start.line; - } - return b.range.start.character - a.range.start.character; - }); - - // Detect overlapping ranges: in reverse-sorted order, each edit's start - // must be >= the next edit's end. If not, the edits would clobber each other - // once applied bottom-up (typically a multi-server rename with stale positions). - for (let i = 0; i < sortedEdits.length - 1; i++) { - const later = sortedEdits[i].range; - const earlier = sortedEdits[i + 1].range; - if (comparePosition(earlier.end, later.start) > 0) { - throw new ToolError( - `overlapping LSP edits: ${formatRange(earlier)} conflicts with ${formatRange(later)}; multi-server rename produced inconsistent edits`, - ); - } - } + const sortedEdits = sortAndValidateTextEdits(edits); for (const edit of sortedEdits) { const { start, end } = edit.range; @@ -78,6 +58,42 @@ export function rangesOverlap(a: Range, b: Range): boolean { return comparePosition(a.start, b.end) < 0 && comparePosition(b.start, a.end) < 0; } +/** + * Sort edits bottom-to-top for in-place application and reject overlaps. + * Equal start positions tiebreak by original array index descending so that, + * applied bottom-up, inserts at the same position land in array order + * (LSP spec: the order of edits in the array defines the order in the result). + */ +export function sortAndValidateTextEdits(edits: TextEdit[]): TextEdit[] { + const sorted = edits + .map((edit, index) => ({ edit, index })) + .sort((a, b) => { + if (a.edit.range.start.line !== b.edit.range.start.line) { + return b.edit.range.start.line - a.edit.range.start.line; + } + if (a.edit.range.start.character !== b.edit.range.start.character) { + return b.edit.range.start.character - a.edit.range.start.character; + } + return b.index - a.index; + }) + .map(entry => entry.edit); + + // Detect overlapping ranges: in reverse-sorted order, each edit's start + // must be >= the next edit's end. If not, the edits would clobber each other + // once applied bottom-up (typically a multi-server rename with stale positions). + for (let i = 0; i < sorted.length - 1; i++) { + const later = sorted[i].range; + const earlier = sorted[i + 1].range; + if (comparePosition(earlier.end, later.start) > 0) { + throw new ToolError( + `overlapping LSP edits: ${formatRange(earlier)} conflicts with ${formatRange(later)}; multi-server rename produced inconsistent edits`, + ); + } + } + + return sorted; +} + /** * Flatten a WorkspaceEdit's text edits into a Map. * Resource operations (create/rename/delete) are ignored — callers handle them separately. @@ -120,92 +136,124 @@ export async function applyTextEdits(filePath: string, edits: TextEdit[]): Promi // Workspace Edit Application // ============================================================================= +type WorkspaceEditOp = + | { kind: "text"; uri: string; edits: TextEdit[] } + | { kind: "create"; uri: string } + | { kind: "rename"; oldUri: string; newUri: string } + | { kind: "delete"; uri: string }; + +/** + * Flatten documentChanges into an ordered op list. Text edits are accumulated + * per-URI and flushed before any resource op that touches the same URI (or, + * for folder rename/delete, any descendant URI) so that renames, creates, and + * deletes always see the correct prior file state. + */ +function planDocumentChanges(documentChanges: NonNullable): WorkspaceEditOp[] { + const ops: WorkspaceEditOp[] = []; + const pending = new Map(); + + const flushUri = (uri: string) => { + const edits = pending.get(uri); + if (!edits) return; + pending.delete(uri); + ops.push({ kind: "text", uri, edits }); + }; + + // Flush the exact URI plus every pending descendant (for folder-level + // resource ops where the queued edits target child files of the target). + const flushSubtree = (uri: string) => { + const prefix = uri.endsWith("/") ? uri : `${uri}/`; + const matches: string[] = []; + for (const candidate of pending.keys()) { + if (candidate === uri || candidate.startsWith(prefix)) matches.push(candidate); + } + for (const target of matches) { + flushUri(target); + } + }; + + for (const change of documentChanges) { + if ("textDocument" in change && change.textDocument && "edits" in change && change.edits) { + const tdc = change as TextDocumentEdit; + const uri = tdc.textDocument.uri; + const textEdits = tdc.edits.filter((e): e is TextEdit => "range" in e && "newText" in e); + if (textEdits.length > 0) { + const prev = pending.get(uri); + if (prev) prev.push(...textEdits); + else pending.set(uri, [...textEdits]); + } + } else if ("kind" in change && change.kind) { + if (change.kind === "create") { + const createOp = change as CreateFile; + flushUri(createOp.uri); + ops.push({ kind: "create", uri: createOp.uri }); + } else if (change.kind === "rename") { + const renameOp = change as RenameFile; + // Per LSP §3.16.2 documentChanges are applied in declared order. + // Flush both the source subtree (so prior edits land before the move) + // AND the destination subtree (so prior edits land on whatever exists + // at newUri before the rename overwrites/replaces it — relevant under + // `options.overwrite` and `options.ignoreIfExists`). + flushSubtree(renameOp.oldUri); + flushSubtree(renameOp.newUri); + ops.push({ kind: "rename", oldUri: renameOp.oldUri, newUri: renameOp.newUri }); + } else if (change.kind === "delete") { + const deleteOp = change as DeleteFile; + flushSubtree(deleteOp.uri); + ops.push({ kind: "delete", uri: deleteOp.uri }); + } + } + } + + // Flush text edits not followed by a resource op. + for (const uri of [...pending.keys()]) { + flushUri(uri); + } + + return ops; +} + /** * Apply a workspace edit (collection of file changes). + * All text-edit batches are overlap-validated before anything is written so a + * conflict throws without leaving the workspace half-applied. * Returns array of applied change descriptions. */ export async function applyWorkspaceEdit(edit: WorkspaceEdit, cwd: string): Promise { const applied: string[] = []; if (edit.documentChanges) { - // Walk documentChanges in original order. Accumulate text edits per-URI and - // flush them before any resource op that touches the same URI (or, for folder - // rename/delete, any descendant URI) so that renames, creates, and deletes - // always see the correct prior file state. - const pending = new Map(); - - const flushUri = async (uri: string) => { - const edits = pending.get(uri); - if (!edits) return; - pending.delete(uri); - const filePath = uriToFile(uri); - await applyTextEdits(filePath, edits); - applied.push(`Applied ${edits.length} edit(s) to ${formatPathRelativeToCwd(filePath, cwd)}`); - }; - - // Flush the exact URI plus every pending descendant (for folder-level - // resource ops where the queued edits target child files of the target). - const flushSubtree = async (uri: string) => { - const prefix = uri.endsWith("/") ? uri : `${uri}/`; - const matches: string[] = []; - for (const candidate of pending.keys()) { - if (candidate === uri || candidate.startsWith(prefix)) matches.push(candidate); - } - for (const target of matches) { - await flushUri(target); - } - }; - - for (const change of edit.documentChanges) { - if ("textDocument" in change && change.textDocument && "edits" in change && change.edits) { - const tdc = change as TextDocumentEdit; - const uri = tdc.textDocument.uri; - const textEdits = tdc.edits.filter((e): e is TextEdit => "range" in e && "newText" in e); - if (textEdits.length > 0) { - const prev = pending.get(uri); - if (prev) prev.push(...textEdits); - else pending.set(uri, [...textEdits]); - } - } else if ("kind" in change && change.kind) { - if (change.kind === "create") { - const createOp = change as CreateFile; - await flushUri(createOp.uri); - const filePath = uriToFile(createOp.uri); - await Bun.write(filePath, ""); - applied.push(`Created ${formatPathRelativeToCwd(filePath, cwd)}`); - } else if (change.kind === "rename") { - const renameOp = change as RenameFile; - // Per LSP §3.16.2 documentChanges are applied in declared order. - // Flush both the source subtree (so prior edits land before the move) - // AND the destination subtree (so prior edits land on whatever exists - // at newUri before the rename overwrites/replaces it — relevant under - // `options.overwrite` and `options.ignoreIfExists`). - await flushSubtree(renameOp.oldUri); - await flushSubtree(renameOp.newUri); - const oldPath = uriToFile(renameOp.oldUri); - const newPath = uriToFile(renameOp.newUri); - await fs.mkdir(path.dirname(newPath), { recursive: true }); - await fs.rename(oldPath, newPath); - applied.push( - `Renamed ${formatPathRelativeToCwd(oldPath, cwd)} → ${formatPathRelativeToCwd(newPath, cwd)}`, - ); - } else if (change.kind === "delete") { - const deleteOp = change as DeleteFile; - await flushSubtree(deleteOp.uri); - const filePath = uriToFile(deleteOp.uri); - await fs.rm(filePath, { recursive: true }); - applied.push(`Deleted ${formatPathRelativeToCwd(filePath, cwd)}`); - } - } + const ops = planDocumentChanges(edit.documentChanges); + for (const op of ops) { + if (op.kind === "text") sortAndValidateTextEdits(op.edits); } - - // Flush text edits not followed by a resource op. - for (const [uri] of pending) { - await flushUri(uri); + for (const op of ops) { + if (op.kind === "text") { + const filePath = uriToFile(op.uri); + await applyTextEdits(filePath, op.edits); + applied.push(`Applied ${op.edits.length} edit(s) to ${formatPathRelativeToCwd(filePath, cwd)}`); + } else if (op.kind === "create") { + const filePath = uriToFile(op.uri); + await Bun.write(filePath, ""); + applied.push(`Created ${formatPathRelativeToCwd(filePath, cwd)}`); + } else if (op.kind === "rename") { + const oldPath = uriToFile(op.oldUri); + const newPath = uriToFile(op.newUri); + await fs.mkdir(path.dirname(newPath), { recursive: true }); + await fs.rename(oldPath, newPath); + applied.push(`Renamed ${formatPathRelativeToCwd(oldPath, cwd)} → ${formatPathRelativeToCwd(newPath, cwd)}`); + } else { + const filePath = uriToFile(op.uri); + await fs.rm(filePath, { recursive: true }); + applied.push(`Deleted ${formatPathRelativeToCwd(filePath, cwd)}`); + } } } else if (edit.changes) { - // Legacy changes-map path: apply all text edits in one pass. + // Legacy changes-map path: validate every file's edits before writing any. const changes = edit.changes; + for (const uri in changes) { + sortAndValidateTextEdits(changes[uri]); + } for (const uri in changes) { const textEdits = changes[uri]; if (textEdits.length === 0) continue; diff --git a/packages/coding-agent/src/lsp/index.ts b/packages/coding-agent/src/lsp/index.ts index 6060dbe95..2fab9f38b 100644 --- a/packages/coding-agent/src/lsp/index.ts +++ b/packages/coding-agent/src/lsp/index.ts @@ -454,21 +454,23 @@ function isMethodNotFoundError(err: unknown): boolean { } async function reloadServer(client: LspClient, serverName: string, signal?: AbortSignal): Promise { - let output = `Restarted ${serverName}`; - const reloadMethods = ["rust-analyzer/reloadWorkspace", "workspace/didChangeConfiguration"]; - for (const method of reloadMethods) { - try { - await sendRequest(client, method, method.includes("Configuration") ? { settings: {} } : null, signal); - output = `Reloaded ${serverName}`; - break; - } catch { - // Method not supported, try next - } + // rust-analyzer exposes a real reload request. + try { + await sendRequest(client, "rust-analyzer/reloadWorkspace", null, signal); + return `Reloaded ${serverName}`; + } catch { + // Method not supported — fall through. } - if (output.startsWith("Restarted")) { + // workspace/didChangeConfiguration is a notification per spec; sending it + // as a request hangs until the tool deadline on servers that route it to + // the notification handler and never respond. + try { + await sendNotification(client, "workspace/didChangeConfiguration", { settings: {} }); + return `Reloaded ${serverName}`; + } catch { client.proc.kill(); + return `Restarted ${serverName}`; } - return output; } interface WaitForDiagnosticsOptions { @@ -636,12 +638,13 @@ interface GetDiagnosticsForFileOptions { async function captureDiagnosticVersions( cwd: string, servers: Array<[string, ServerConfig]>, + initTimeoutMs?: number, ): Promise { const versions = new Map(); await Promise.allSettled( servers.map(async ([serverName, serverConfig]) => { if (serverConfig.createClient) return; - const client = await getOrCreateClient(serverConfig, cwd); + const client = await getOrCreateClient(serverConfig, cwd, initTimeoutMs); versions.set(serverName, client.diagnosticsVersion); }), ); @@ -1118,7 +1121,9 @@ async function runLspWritethrough( const useCustomFormatter = enableFormat && customLinterServers.length > 0; // Capture diagnostic versions BEFORE syncing to detect stale diagnostics - const minVersions = enableDiagnostics ? await captureDiagnosticVersions(cwd, servers) : undefined; + // Bound client creation by the writethrough budget: a hung/broken server + // must not add its full init wait (30s default) to every edit. + const minVersions = enableDiagnostics ? await captureDiagnosticVersions(cwd, servers, 5_000) : undefined; let expectedDocumentVersions: ServerVersionMap | undefined; let formatter: FileFormatResult | undefined; @@ -2311,11 +2316,12 @@ export class LspTool implements AgentTool - (parsedIndex !== null && index === parsedIndex) || - actionItem.title.toLowerCase().includes(normalizedQuery.toLowerCase()), - ); + const selectedAction = + parsedIndex !== null + ? result[parsedIndex] + : result.find(actionItem => + actionItem.title.toLowerCase().includes(normalizedQuery.toLowerCase()), + ); if (!selectedAction) { const actionLines = result.map((actionItem, index) => ` ${formatCodeAction(actionItem, index)}`); diff --git a/packages/coding-agent/src/lsp/types.ts b/packages/coding-agent/src/lsp/types.ts index 42028047a..84346e5f4 100644 --- a/packages/coding-agent/src/lsp/types.ts +++ b/packages/coding-agent/src/lsp/types.ts @@ -416,6 +416,8 @@ export interface LspClient { pendingRequests: Map; messageBuffer: Uint8Array; isReading: boolean; + /** Lifecycle state: "connecting" until initialize completes, then "ready"; "error" on init failure or reader death. */ + status: "connecting" | "ready" | "error"; serverCapabilities?: LspServerCapabilities; lastActivity: number; /** Serializes outbound JSON-RPC writes to the server process. */ diff --git a/packages/coding-agent/src/lsp/utils.ts b/packages/coding-agent/src/lsp/utils.ts index 768d706c2..d9eeac7ac 100644 --- a/packages/coding-agent/src/lsp/utils.ts +++ b/packages/coding-agent/src/lsp/utils.ts @@ -27,22 +27,17 @@ export { detectLanguageId } from "../utils/lang-from-path"; /** * Convert a file path to a file:// URI. + * Uses the URL machinery so special characters (`%`, `#`, `?`, spaces) are + * percent-encoded; plain concatenation produced URIs that broke round-trips. * Handles Windows drive letters correctly. */ export function fileToUri(filePath: string): string { - const resolved = path.resolve(filePath); - - if (process.platform === "win32") { - // Windows: file:///C:/path/to/file - return `file:///${resolved.replace(/\\/g, "/")}`; - } - - // Unix: file:///path/to/file - return `file://${resolved}`; + return Bun.pathToFileURL(path.resolve(filePath)).href; } /** * Convert a file:// URI to a file path. + * Tolerates both percent-encoded URIs and lax servers that send raw paths. * Handles Windows drive letters correctly. */ export function uriToFile(uri: string): string { @@ -50,7 +45,30 @@ export function uriToFile(uri: string): string { return uri; } - let filePath = decodeURIComponent(uri.slice(7)); + // A raw `#`/`?` parses *successfully* as fragment/query and silently + // truncates the path — it never reaches the catch below. LSP servers do + // not use fragments or queries on file URIs (encoded forms are %23/%3F), + // so raw occurrences mean a lax server sent an unencoded path. + if (uri.includes("#") || uri.includes("?")) { + return laxUriToFile(uri); + } + + try { + return Bun.fileURLToPath(uri); + } catch { + // Not a well-formed file URL (unencoded characters, stray `%`, host + // component). Fall back to a lenient manual conversion. + return laxUriToFile(uri); + } +} + +function laxUriToFile(uri: string): string { + let filePath = uri.slice(7); + try { + filePath = decodeURIComponent(filePath); + } catch { + // Invalid percent-encoding — treat as a literal path. + } // Windows: file:///C:/path → C:/path (strip leading slash before drive letter) if (process.platform === "win32" && filePath.startsWith("/") && /^[A-Za-z]:/.test(filePath.slice(1))) { diff --git a/packages/coding-agent/test/tools/lsp-diagnostics-freshness.test.ts b/packages/coding-agent/test/tools/lsp-diagnostics-freshness.test.ts index de79838b8..81e55ce3e 100644 --- a/packages/coding-agent/test/tools/lsp-diagnostics-freshness.test.ts +++ b/packages/coding-agent/test/tools/lsp-diagnostics-freshness.test.ts @@ -37,6 +37,7 @@ function createClient(cwd: string, config: ServerConfig): LspClient { pendingRequests: new Map(), messageBuffer: new Uint8Array(), isReading: false, + status: "ready", lastActivity: Date.now(), writeQueue: Promise.resolve(), activeProgressTokens: new Set(), diff --git a/packages/coding-agent/test/tools/lsp-regressions.test.ts b/packages/coding-agent/test/tools/lsp-regressions.test.ts index ca5e81765..4b323fdc3 100644 --- a/packages/coding-agent/test/tools/lsp-regressions.test.ts +++ b/packages/coding-agent/test/tools/lsp-regressions.test.ts @@ -7,7 +7,7 @@ import { LspTool } from "@oh-my-pi/pi-coding-agent/lsp"; import * as lspClient from "@oh-my-pi/pi-coding-agent/lsp/client"; import * as lspConfig from "@oh-my-pi/pi-coding-agent/lsp/config"; import { getServersForFile, loadConfig } from "@oh-my-pi/pi-coding-agent/lsp/config"; -import { applyWorkspaceEdit } from "@oh-my-pi/pi-coding-agent/lsp/edits"; +import { applyTextEditsToString, applyWorkspaceEdit } from "@oh-my-pi/pi-coding-agent/lsp/edits"; import { renderCall, renderResult } from "@oh-my-pi/pi-coding-agent/lsp/render"; import type { CodeAction, @@ -31,6 +31,7 @@ import { hasGlobPattern, resolveDiagnosticTargets, resolveSymbolColumn, + uriToFile, } from "@oh-my-pi/pi-coding-agent/lsp/utils"; import { getThemeByName } from "@oh-my-pi/pi-coding-agent/modes/theme/theme"; import type { ToolSession } from "@oh-my-pi/pi-coding-agent/tools"; @@ -799,6 +800,7 @@ for await (const chunk of Bun.stdin.stream()) { pendingRequests: new Map(), messageBuffer: new Uint8Array(), isReading: false, + status: "ready", lastActivity: Date.now(), writeQueue: Promise.resolve(), activeProgressTokens: new Set(), @@ -1002,6 +1004,7 @@ for await (const chunk of Bun.stdin.stream()) { pendingRequests: new Map(), messageBuffer: new Uint8Array(), isReading: false, + status: "ready", lastActivity: Date.now(), writeQueue: Promise.resolve(), activeProgressTokens: new Set(), @@ -1103,6 +1106,7 @@ for await (const chunk of Bun.stdin.stream()) { pendingRequests: new Map(), messageBuffer: new Uint8Array(), isReading: false, + status: "ready", lastActivity: Date.now(), writeQueue: Promise.resolve(), activeProgressTokens: new Set(), @@ -1163,6 +1167,7 @@ for await (const chunk of Bun.stdin.stream()) { pendingRequests: new Map(), messageBuffer: new Uint8Array(), isReading: false, + status: "ready", lastActivity: Date.now(), writeQueue: Promise.resolve(), activeProgressTokens: new Set(), @@ -1232,6 +1237,7 @@ for await (const chunk of Bun.stdin.stream()) { pendingRequests: new Map(), messageBuffer: new Uint8Array(), isReading: false, + status: "ready", lastActivity: Date.now(), writeQueue: Promise.resolve(), activeProgressTokens: new Set(), @@ -1301,6 +1307,7 @@ for await (const chunk of Bun.stdin.stream()) { pendingRequests: new Map(), messageBuffer: new Uint8Array(), isReading: false, + status: "ready", lastActivity: Date.now(), writeQueue: Promise.resolve(), activeProgressTokens: new Set(), @@ -1358,6 +1365,7 @@ for await (const chunk of Bun.stdin.stream()) { pendingRequests: new Map(), messageBuffer: new Uint8Array(), isReading: false, + status: "ready", lastActivity: Date.now(), writeQueue: Promise.resolve(), activeProgressTokens: new Set(), @@ -1506,6 +1514,67 @@ for await (const chunk of Bun.stdin.stream()) { tempDir.removeSync(); } }); + + it("applies equal-position inserts in array order", () => { + // LSP spec: multiple inserts at the same position land in the order they + // appear in the edits array (import + reference insertions rely on this). + const result = applyTextEditsToString("abc", [ + { range: { start: { line: 0, character: 1 }, end: { line: 0, character: 1 } }, newText: "X" }, + { range: { start: { line: 0, character: 1 }, end: { line: 0, character: 1 } }, newText: "Y" }, + ]); + expect(result).toBe("aXYbc"); + }); + + it("validates every file's edits before writing any workspace-edit file", async () => { + const tempDir = TempDir.createSync("@omp-lsp-atomic-validate-"); + try { + const okPath = path.join(tempDir.path(), "ok.ts"); + const badPath = path.join(tempDir.path(), "bad.ts"); + const okContent = "export const ok = 1;\n"; + await Bun.write(okPath, okContent); + await Bun.write(badPath, "export const bad = 2;\n"); + + const workspaceEdit: WorkspaceEdit = { + changes: { + [fileToUri(okPath)]: [ + { + range: { start: { line: 0, character: 13 }, end: { line: 0, character: 15 } }, + newText: "changed", + }, + ], + [fileToUri(badPath)]: [ + // Overlapping edits — must reject the whole workspace edit. + { + range: { start: { line: 0, character: 0 }, end: { line: 0, character: 10 } }, + newText: "x", + }, + { + range: { start: { line: 0, character: 5 }, end: { line: 0, character: 12 } }, + newText: "y", + }, + ], + }, + }; + + await expect(applyWorkspaceEdit(workspaceEdit, tempDir.path())).rejects.toThrow(/overlapping LSP edits/); + // The valid file must be untouched: validation runs before any write. + expect(fs.readFileSync(okPath, "utf8")).toBe(okContent); + } finally { + tempDir.removeSync(); + } + }); + + it("round-trips file URIs containing percent and hash characters", () => { + const tricky = path.join("/tmp", "omp uri", "100% #1.ts"); + const uri = fileToUri(tricky); + // Percent-encoded so the server cannot misparse a fragment or escape. + expect(uri).not.toContain("#"); + expect(uri).not.toContain(" "); + expect(uriToFile(uri)).toBe(tricky); + // Lax servers sending unencoded paths are tolerated. + expect(uriToFile("file:///tmp/omp uri/plain.ts")).toBe("/tmp/omp uri/plain.ts"); + }); + it("resolves $-prefixed identifiers past compound matches", async () => { // Pre-fix, BARE_IDENTIFIER_RE rejected leading `$`, so requireWordBoundary // was false and `resolveSymbolColumn(_, _, "$store")` returned the column @@ -1645,6 +1714,7 @@ for await (const chunk of Bun.stdin.stream()) { pendingRequests: new Map(), messageBuffer: new Uint8Array(), isReading: false, + status: "ready", lastActivity: Date.now(), writeQueue: Promise.resolve(), activeProgressTokens: new Set(), @@ -1670,6 +1740,7 @@ for await (const chunk of Bun.stdin.stream()) { pendingRequests: new Map(), messageBuffer: new Uint8Array(), isReading: false, + status: "ready", lastActivity: Date.now(), writeQueue: Promise.resolve(), activeProgressTokens: new Set(),