Merge PR #6719: fix(coding-agent): handle asynchronous LSP pipe failures (@shoucandanghehe)

This commit is contained in:
can1357
2026-07-27 04:58:27 +02:00
3 changed files with 65 additions and 18 deletions
+1
View File
@@ -7,6 +7,7 @@
- Fixed `glob` rejecting safe `memory://root/<directory>/**` patterns. Memory globs now resolve their directory prefix inside the project memory root while rejecting traversal and percent-encoded path separators across the complete glob path.
- Fixed `omp --resume <id>` prompting to fork sessions from another existing directory instead of switching the process and cwd-scoped settings into the resumed session's recorded directory ([#6752](https://github.com/can1357/oh-my-pi/issues/6752)).
- Fixed deferred CLI model roles resolving ambiguous bare model IDs to a preferred but unauthenticated provider instead of the authenticated provider selected by the eager path ([#6727](https://github.com/can1357/oh-my-pi/issues/6727)).
- Fixed Windows sessions crashing with an unhandled `EPIPE: broken pipe, write` when an LSP server closed its stdin between filesystem mutations; LSP writes now observe asynchronous `FileSink.write()` failures and route them through the existing request/notification failure path.
## [17.1.4] - 2026-07-26
+22 -18
View File
@@ -200,10 +200,10 @@ function abortReason(signal: AbortSignal): Error {
return signal.reason instanceof Error ? signal.reason : new ToolAbortError();
}
class LspFlushAbortError extends Error {
class LspDrainAbortError extends Error {
constructor(readonly reason: Error) {
super(reason.message);
this.name = "LspFlushAbortError";
this.name = "LspDrainAbortError";
}
}
@@ -216,26 +216,30 @@ async function writeMessage(
throw abortReason(signal);
}
const content = JSON.stringify(message);
sink.write(`Content-Length: ${Buffer.byteLength(content, "utf-8")}\r\n\r\n${content}`);
const flush = Promise.resolve(sink.flush());
const write = Promise.resolve(
sink.write(`Content-Length: ${Buffer.byteLength(content, "utf-8")}\r\n\r\n${content}`),
);
// Attach before flush(): it may throw synchronously after write() returned a
// rejected Promise, and leaving that rejection unobserved kills the host.
void write.catch(() => {});
const drain = Promise.all([write, Promise.resolve(sink.flush())]).then(() => {});
if (!signal) {
await flush;
await drain;
return;
}
// The sink's flush blocks on the OS-level pipe drain: if the server is
// alive but stopped reading stdin, `await sink.flush()` never resolves.
// Race the flush against the caller's signal so a wedged server surfaces
// as the tool's normal timeout/cancel instead of a permanent hang.
// Either sink operation can block on the OS-level pipe drain when a live
// server stops reading stdin. Race the combined drain against the caller's
// signal so a wedged server surfaces as the tool's normal timeout/cancel.
const { promise, resolve, reject } = Promise.withResolvers<void>();
const onAbort = () => {
signal.removeEventListener("abort", onAbort);
// The underlying flush stays pending in the background; suppress its
// The underlying drain stays pending in the background; suppress its
// eventual settlement so we do not surface an unhandled rejection.
flush.catch(() => {});
reject(new LspFlushAbortError(abortReason(signal)));
drain.catch(() => {});
reject(new LspDrainAbortError(abortReason(signal)));
};
signal.addEventListener("abort", onAbort, { once: true });
flush.then(
drain.then(
() => {
signal.removeEventListener("abort", onAbort);
resolve();
@@ -249,8 +253,8 @@ async function writeMessage(
}
/**
* Kill a client whose write queue is stuck (aborted flush left the sink's
* flush promise pending, so subsequent writes queue behind a wedge forever).
* Kill a client whose write queue is stuck (an aborted drain left a sink
* operation pending, so subsequent writes queue behind the wedge forever).
* Remove it from `clients` immediately so concurrent `getOrCreateClient`
* callers do not grab the corpse before `proc.exited` cleans up.
*/
@@ -270,8 +274,8 @@ function queueWriteMessage(
): Promise<void> {
const write = client.writeQueue.catch(() => {}).then(() => writeMessage(client.proc.stdin, message, signal));
const result = write.catch((err: unknown) => {
if (err instanceof LspFlushAbortError) {
// Only an abort that raced this write's in-flight flush leaves
if (err instanceof LspDrainAbortError) {
// Only an abort that raced this write's in-flight drain leaves
// the sink pending. Pre-write aborts and queued caller timeouts
// must not kill a healthy shared client.
teardownWedgedClient(client);
@@ -1067,7 +1071,7 @@ const WATCHED_FILES_NOTIFY_TIMEOUT_MS = 2_000;
* This covers sibling files that are not open text documents, such as generated
* CSS modules or type files that another edited document imports immediately.
*
* The underlying stdin flush is self-bounded by
* The underlying stdin write drain is self-bounded by
* {@link WATCHED_FILES_NOTIFY_TIMEOUT_MS}; only an abort of the caller's
* `signal` rejects.
*/
@@ -3079,6 +3079,48 @@ describe("lsp regressions", () => {
expect(writes).toHaveLength(1);
});
it("surfaces an asynchronous stdin write rejection instead of resolving the notification", async () => {
const epipe = Object.assign(new Error("EPIPE: broken pipe, write"), {
code: "EPIPE",
syscall: "write",
});
const client: LspClient = {
name: "fake-lsp-write-epipe:/tmp",
cwd: "/tmp",
config: { command: "fake-lsp-write-epipe", fileTypes: ["ts"], rootMarkers: [] },
proc: {
exited: Promise.withResolvers<number>().promise,
exitCode: null,
stdin: {
write: () => Promise.reject(epipe),
flush: () => 0,
},
stdout: new ReadableStream<Uint8Array>(),
peekStderr: () => "",
kill() {},
} as unknown as LspClient["proc"],
requestId: 0,
diagnostics: new Map(),
diagnosticsVersion: 0,
openFiles: new Map(),
pendingRequests: new Map(),
messageBuffer: new Uint8Array(0),
isReading: false,
status: "ready",
lastActivity: Date.now(),
writeQueue: Promise.resolve(),
activeProgressTokens: new Set(),
projectLoaded: Promise.resolve(),
resolveProjectLoaded: () => {},
};
await expect(
lspClient.sendNotification(client, "workspace/didChangeWatchedFiles", {
changes: [{ uri: "file:///tmp/example.ts", type: 2 }],
}),
).rejects.toBe(epipe);
});
it("bounds a wedged notification flush on the caller signal and tears down the client", async () => {
// Custom fake: stdin.flush is gated by a controllable promise so we
// can simulate a server that stopped draining stdin AFTER init has