refactor(coding-agent): reduced lsp index to a composition barrel

- src/lsp/index.ts is the explicit ./lsp package entry, yet held 2821 lines of
  warmup, config caching, diagnostics, external build-command workspace
  diagnostics, the writethrough batching subsystem and the LspTool class.
- Those are now servers, diagnostics, workspace-diagnostics, writethrough and
  tool modules; index.ts is 22 lines and re-exports the same public surface.
- configCache and writethroughBatches remain single instances and every tuned
  diagnostics timing constant moved verbatim.
This commit is contained in:
can1357
2026-08-08 06:32:01 +02:00
parent 0f3e45f07a
commit cd4e04e8ee
6 changed files with 2915 additions and 2819 deletions
@@ -0,0 +1,516 @@
import * as fs from "node:fs";
import path from "node:path";
import { logger } from "@oh-my-pi/pi-utils";
import { formatPathRelativeToCwd } from "../tools/path-utils";
import { throwIfAborted } from "../tools/tool-errors";
import { getOrCreateClient, sendRequest, supportsDocumentDiagnostics, waitForProjectLoaded } from "./client";
import { getLinterClient } from "./clients";
import { hasRootMarkerAncestor } from "./config";
import { applyTextEditsToString } from "./edits";
import { resolveFormatOptions } from "./format-options";
import { isProjectAwareLspServer } from "./servers";
import type {
Diagnostic,
Location,
LocationLink,
LspClient,
Position,
PublishedDiagnostics,
ServerConfig,
TextEdit,
} from "./types";
import {
fileToUri,
formatDiagnostic,
formatDiagnosticsSummary,
formatLocation,
readLocationContext,
sortDiagnostics,
uriToFile,
} from "./utils";
const DIAGNOSTIC_MESSAGE_LIMIT = 50;
export const SINGLE_DIAGNOSTICS_WAIT_TIMEOUT_MS = 3000;
export const BATCH_DIAGNOSTICS_WAIT_TIMEOUT_MS = 400;
const DIAGNOSTICS_POLL_MS = 100;
const DIAGNOSTICS_SETTLE_MS = 250;
/**
* How long the edit/write writethrough blocks inline waiting for fresh
* diagnostics before handing slow servers off to the deferred late-injection
* channel. Keeps the common fast-server case inline while letting an edit
* return promptly when a server (e.g. a large-monorepo tsserver) is slow to
* publish fresh diagnostics.
*/
export const INLINE_DIAGNOSTICS_WAIT_TIMEOUT_MS = 500;
/**
* Inner per-server diagnostics wait budget for the background/deferred fetch.
* Longer than the inline cap (and the old 3s default) so a slow server still
* delivers late instead of giving up before it ever publishes.
*/
export const DEFERRED_DIAGNOSTICS_WAIT_TIMEOUT_MS = 12_000;
export const MAX_GLOB_DIAGNOSTIC_TARGETS = 20;
export const WORKSPACE_SYMBOL_LIMIT = 200;
export const PROJECT_INDEXED_ACTIONS: ReadonlySet<string> = new Set([
"definition",
"type_definition",
"implementation",
"references",
"rename",
"hover",
]);
const RUST_WORKSPACE_MARKERS = ["Cargo.toml", "rust-analyzer.toml"] as const;
export function hasRustWorkspaceAncestor(filePath: string): boolean {
let dir = path.dirname(filePath);
while (true) {
for (const marker of RUST_WORKSPACE_MARKERS) {
if (fs.existsSync(path.join(dir, marker))) {
return true;
}
}
const parent = path.dirname(dir);
if (parent === dir) {
return false;
}
dir = parent;
}
}
export function limitDiagnosticMessages(messages: string[]): string[] {
if (messages.length <= DIAGNOSTIC_MESSAGE_LIMIT) {
return messages;
}
return messages.slice(0, DIAGNOSTIC_MESSAGE_LIMIT);
}
const ORPHAN_TYPESCRIPT_PROJECT_DIAGNOSTIC_CODES: Record<number, true> = {
1375: true,
1378: true,
2307: true,
2580: true,
2591: true,
2792: true,
2867: true,
};
function diagnosticCodeNumber(diagnostic: Diagnostic): number | null {
if (typeof diagnostic.code === "number") return diagnostic.code;
if (typeof diagnostic.code === "string" && /^\d+$/.test(diagnostic.code)) return Number(diagnostic.code);
return null;
}
function isTypeScriptProjectDiagnostic(serverName: string, diagnostic: Diagnostic): boolean {
if (diagnostic.source !== "typescript" && !serverName.toLowerCase().includes("typescript")) {
return false;
}
const code = diagnosticCodeNumber(diagnostic);
return code !== null && ORPHAN_TYPESCRIPT_PROJECT_DIAGNOSTIC_CODES[code] === true;
}
function filterOrphanProjectDiagnostics(
absolutePath: string,
serverName: string,
serverConfig: ServerConfig,
diagnostics: Diagnostic[],
): Diagnostic[] {
if (!serverConfig.rootMarkers.length || hasRootMarkerAncestor(absolutePath, serverConfig.rootMarkers)) {
return diagnostics;
}
return diagnostics.filter(diagnostic => !isTypeScriptProjectDiagnostic(serverName, diagnostic));
}
const LOCATION_CONTEXT_LINES = 1;
export const REFERENCE_CONTEXT_LIMIT = 50;
export const REFERENCES_RETRY_COUNT = 2;
export const REFERENCES_RETRY_DELAY_MS = 250;
function comparePosition(a: Position, b: Position): number {
return a.line === b.line ? a.character - b.character : a.line - b.line;
}
function rangeContainsPosition(range: Location["range"], position: Position): boolean {
return comparePosition(range.start, position) <= 0 && comparePosition(position, range.end) <= 0;
}
export function isOnlyQueriedDeclaration(locations: Location[], uri: string, position: Position): boolean {
return locations.length === 1 && locations[0]?.uri === uri && rangeContainsPosition(locations[0].range, position);
}
export function normalizeLocationResult(
result: Location | Location[] | LocationLink | LocationLink[] | null,
): Location[] {
if (!result) return [];
const raw = Array.isArray(result) ? result : [result];
return raw.flatMap(loc => {
if ("uri" in loc) {
return [loc as Location];
}
if ("targetUri" in loc) {
const link = loc as LocationLink;
return [{ uri: link.targetUri, range: link.targetSelectionRange ?? link.targetRange }];
}
return [];
});
}
export async function formatLocationWithContext(location: Location, cwd: string): Promise<string> {
const header = ` ${formatLocation(location, cwd)}`;
const context = await readLocationContext(
uriToFile(location.uri),
location.range.start.line + 1,
LOCATION_CONTEXT_LINES,
);
if (context.length === 0) {
return header;
}
return `${header}\n${context.map(lineText => ` ${lineText}`).join("\n")}`;
}
interface WaitForDiagnosticsOptions {
timeoutMs?: number;
signal?: AbortSignal;
minVersion?: number;
expectedDocumentVersion?: number;
/**
* Quiescence window (ms). typescript-language-server never echoes the document
* version (issue #983) and emits diagnostics from several sources at different
* times, so there is no single "complete, version-matched" publish to gate on.
* When the server does not exact-version-match, accept the latest publish only
* after no newer one has arrived for this long, letting an in-flight pre-edit
* publish be superseded by the fresh one.
*/
settleMs?: number;
}
function requestDocumentDiagnostics(
client: LspClient,
uri: string,
signal: AbortSignal | undefined,
timeoutMs: number,
): Promise<Diagnostic[] | undefined> {
return sendRequest(client, "textDocument/diagnostic", { textDocument: { uri } }, signal, timeoutMs)
.then(report => {
if (!report || typeof report !== "object" || !("kind" in report) || report.kind !== "full") {
return undefined;
}
if (!("items" in report) || !Array.isArray(report.items)) return undefined;
return report.items;
})
.catch(err => {
if (!signal?.aborted) {
logger.debug("LSP document diagnostic pull failed", { server: client.name, uri, error: String(err) });
}
return undefined;
});
}
export async function waitForDiagnostics(
client: LspClient,
uri: string,
options: WaitForDiagnosticsOptions = {},
): Promise<Diagnostic[]> {
const { timeoutMs = 3000, signal, minVersion, expectedDocumentVersion, settleMs = DIAGNOSTICS_SETTLE_MS } = options;
const deadline = Date.now() + timeoutMs;
let pullAttempted = false;
let pullResultPromise: Promise<{ diagnostics: Diagnostic[] | undefined }> | undefined;
let pulled: Diagnostic[] | undefined;
let settledRef: PublishedDiagnostics | undefined;
let settledAt = 0;
while (Date.now() < deadline) {
throwIfAborted(signal);
if (!pullAttempted && supportsDocumentDiagnostics(client)) {
pullAttempted = true;
pullResultPromise = requestDocumentDiagnostics(client, uri, signal, Math.max(1, deadline - Date.now())).then(
diagnostics => ({ diagnostics }),
);
}
const versionOk = minVersion === undefined || client.diagnosticsVersion > minVersion;
const published = client.diagnostics.get(uri);
if (published && versionOk) {
// Server honored our exact document version → authoritative, accept now.
if (expectedDocumentVersion !== undefined && published.version === expectedDocumentVersion) {
return published.diagnostics;
}
// Unversioned/mismatched publish: wait for the stream to go quiet so an
// in-flight publish for the pre-edit content is superseded by the fresh one.
if (published !== settledRef) {
settledRef = published;
settledAt = Date.now();
} else if (Date.now() - settledAt >= settleMs) {
return published.diagnostics;
}
}
const pollMs = Math.min(DIAGNOSTICS_POLL_MS, Math.max(0, deadline - Date.now()));
if (!pullResultPromise) {
await Bun.sleep(pollMs);
continue;
}
const pullResult = await Promise.race([pullResultPromise, Bun.sleep(pollMs).then(() => undefined)]);
if (pullResult) {
pullResultPromise = undefined;
pulled = pullResult.diagnostics;
if (pulled !== undefined) break;
}
}
const versionOk = minVersion === undefined || client.diagnosticsVersion > minVersion;
const published = client.diagnostics.get(uri);
if (published && versionOk) {
return published.diagnostics;
}
if (pullResultPromise) {
pulled = (await pullResultPromise).diagnostics;
}
throwIfAborted(signal);
if (pulled === undefined) return [];
client.diagnostics.set(uri, {
diagnostics: pulled,
version: expectedDocumentVersion ?? client.openFiles.get(uri)?.version ?? null,
});
client.diagnosticsVersion += 1;
return pulled;
}
/** Result from getDiagnosticsForFile */
export interface FileDiagnosticsResult {
/** Name of the LSP server used (if available) */
server?: string;
/** Formatted diagnostic messages */
messages: string[];
/** Summary string (e.g., "2 error(s), 1 warning(s)") */
summary: string;
/** Whether there are any errors (severity 1) */
errored: boolean;
/** Whether the file was formatted */
formatter?: FileFormatResult;
}
export type ServerVersionMap = Map<string, number>;
interface GetDiagnosticsForFileOptions {
signal?: AbortSignal;
minVersions?: ServerVersionMap;
expectedDocumentVersions?: ServerVersionMap;
/** Per-server wait budget (ms). Defaults to {@link SINGLE_DIAGNOSTICS_WAIT_TIMEOUT_MS}. */
timeoutMs?: number;
}
/**
* Capture current diagnostic versions for all LSP servers.
* Call this BEFORE syncing content to detect stale diagnostics later.
*/
export async function captureDiagnosticVersions(
cwd: string,
servers: Array<[string, ServerConfig]>,
initTimeoutMs?: number,
signal?: AbortSignal,
): Promise<ServerVersionMap> {
const versions = new Map<string, number>();
await Promise.allSettled(
servers.map(async ([serverName, serverConfig]) => {
if (serverConfig.createClient) return;
const client = await getOrCreateClient(serverConfig, cwd, initTimeoutMs, signal);
versions.set(serverName, client.diagnosticsVersion);
}),
);
return versions;
}
export async function captureOpenFileVersions(
absolutePath: string,
cwd: string,
servers: Array<[string, ServerConfig]>,
signal?: AbortSignal,
): Promise<ServerVersionMap> {
const uri = fileToUri(absolutePath);
const versions = new Map<string, number>();
await Promise.allSettled(
servers.map(async ([serverName, serverConfig]) => {
const client = await getOrCreateClient(serverConfig, cwd, undefined, signal);
const version = client.openFiles.get(uri)?.version;
if (version !== undefined) {
versions.set(serverName, version);
}
}),
);
return versions;
}
/**
* Get diagnostics for a file using LSP or custom linter client.
*
* @param absolutePath - Absolute path to the file
* @param cwd - Working directory for LSP config resolution
* @param servers - Servers to query diagnostics for
* @param minVersions - Minimum diagnostic versions per server (to detect stale results)
* @returns Diagnostic results or undefined if no servers
*/
export async function getDiagnosticsForFile(
absolutePath: string,
cwd: string,
servers: Array<[string, ServerConfig]>,
options: GetDiagnosticsForFileOptions = {},
): Promise<FileDiagnosticsResult | undefined> {
const { signal, minVersions, expectedDocumentVersions, timeoutMs } = options;
if (servers.length === 0) {
return undefined;
}
const uri = fileToUri(absolutePath);
const relPath = formatPathRelativeToCwd(absolutePath, cwd);
const allDiagnostics: Diagnostic[] = [];
const serverNames: string[] = [];
// Wait for diagnostics from all servers in parallel
const results = await Promise.allSettled(
servers.map(async ([serverName, serverConfig]) => {
throwIfAborted(signal);
// Use custom linter client if configured
if (serverConfig.createClient) {
const linterClient = getLinterClient(serverName, serverConfig, cwd);
const diagnostics = await linterClient.lint(absolutePath);
return { serverName, serverConfig, diagnostics };
}
// Default: use LSP
const client = await getOrCreateClient(serverConfig, cwd, undefined, signal);
throwIfAborted(signal);
if (isProjectAwareLspServer(serverConfig)) {
await waitForProjectLoaded(client, signal);
throwIfAborted(signal);
}
// Content already synced + didSave sent, wait for fresh diagnostics
const minVersion = minVersions?.get(serverName);
const expectedDocumentVersion = expectedDocumentVersions?.get(serverName);
const diagnostics = await waitForDiagnostics(client, uri, {
timeoutMs: timeoutMs ?? SINGLE_DIAGNOSTICS_WAIT_TIMEOUT_MS,
signal,
minVersion,
expectedDocumentVersion,
});
return { serverName, serverConfig, diagnostics };
}),
);
for (const result of results) {
if (result.status === "fulfilled") {
serverNames.push(result.value.serverName);
allDiagnostics.push(
...filterOrphanProjectDiagnostics(
absolutePath,
result.value.serverName,
result.value.serverConfig,
result.value.diagnostics,
),
);
}
}
if (serverNames.length === 0) {
return undefined;
}
if (allDiagnostics.length === 0) {
return {
server: serverNames.join(", "),
messages: [],
summary: "OK",
errored: false,
};
}
// Deduplicate diagnostics by range + message (different servers might report similar issues)
const seen = new Set<string>();
const uniqueDiagnostics: Diagnostic[] = [];
for (const d of allDiagnostics) {
const key = `${d.range.start.line}:${d.range.start.character}:${d.range.end.line}:${d.range.end.character}:${d.message}`;
if (!seen.has(key)) {
seen.add(key);
uniqueDiagnostics.push(d);
}
}
sortDiagnostics(uniqueDiagnostics);
const formatted = uniqueDiagnostics.map(d => formatDiagnostic(d, relPath));
const limited = limitDiagnosticMessages(formatted);
const summary = formatDiagnosticsSummary(uniqueDiagnostics);
const hasErrors = uniqueDiagnostics.some(d => d.severity === 1);
return {
server: serverNames.join(", "),
messages: limited,
summary,
errored: hasErrors,
};
}
export enum FileFormatResult {
UNCHANGED = "unchanged",
FORMATTED = "formatted",
}
/**
* Format content using LSP or custom linter client.
*
* @param absolutePath - Absolute path (for URI)
* @param content - Content to format
* @param cwd - Working directory for LSP config resolution
* @param servers - Servers to try formatting with
* @returns Formatted content, or original if no formatter available
*/
export async function formatContent(
absolutePath: string,
content: string,
cwd: string,
servers: Array<[string, ServerConfig]>,
signal?: AbortSignal,
): Promise<string> {
if (servers.length === 0) {
return content;
}
const uri = fileToUri(absolutePath);
for (const [serverName, serverConfig] of servers) {
try {
throwIfAborted(signal);
// Use custom linter client if configured
if (serverConfig.createClient) {
const linterClient = getLinterClient(serverName, serverConfig, cwd);
return await linterClient.format(absolutePath, content);
}
// Default: use LSP
const client = await getOrCreateClient(serverConfig, cwd, undefined, signal);
throwIfAborted(signal);
const caps = client.serverCapabilities;
if (!caps?.documentFormattingProvider) {
continue;
}
// Request formatting (content already synced)
const edits = (await sendRequest(
client,
"textDocument/formatting",
{
textDocument: { uri },
options: resolveFormatOptions(absolutePath, content),
},
signal,
)) as TextEdit[] | null;
if (!edits || edits.length === 0) {
return content;
}
// Apply edits in-memory and return
return applyTextEditsToString(content, edits);
} catch {}
}
return content;
}
File diff suppressed because it is too large Load Diff
+296
View File
@@ -0,0 +1,296 @@
import { logger } from "@oh-my-pi/pi-utils";
import { throwIfAborted } from "../tools/tool-errors";
import {
getActiveClients,
getActiveOrPendingClient,
getOrCreateClient,
type LspServerStatus,
notifySaved,
sendNotification,
sendRequest,
setIdleTimeout,
shutdownClientInstance,
syncContent,
WARMUP_TIMEOUT_MS,
} from "./client";
import { getServersForFile, type LspConfig, loadConfig } from "./config";
import { MUX_RESTART_METHOD } from "./mux/protocol";
import type { LspClient, ServerConfig } from "./types";
/**
* LSP actions that do not mutate the workspace or language-server state.
* Anything not in this set (rename, code_actions with apply, rename_file,
* reload, raw request, etc.) is classified as write-tier.
*/
export const LSP_READONLY_ACTIONS: ReadonlySet<string> = new Set([
"diagnostics",
"definition",
"type_definition",
"implementation",
"references",
"hover",
"symbols",
"status",
"capabilities",
]);
export interface LspStartupServerInfo {
name: string;
status: "connecting" | "ready" | "error" | "available";
fileTypes: string[];
error?: string;
}
/** Result from warming up LSP servers */
export interface LspWarmupResult {
servers: Array<LspStartupServerInfo & { status: "ready" | "error" }>;
}
/** Options for warming up LSP servers */
export interface LspWarmupOptions {
/** Called when starting to connect to servers */
onConnecting?: (serverNames: string[]) => void;
}
export function discoverStartupLspServers(
cwd: string,
status: LspStartupServerInfo["status"] = "connecting",
): LspStartupServerInfo[] {
const config = loadConfig(cwd);
return getLspServers(config).map(([name, serverConfig]) => ({
name,
status,
fileTypes: serverConfig.fileTypes,
}));
}
/**
* Warm up LSP servers for a directory by connecting to all detected servers.
* This should be called at startup to avoid cold-start delays.
*
* @param cwd - Working directory to detect and start servers for
* @param options - Optional callbacks for progress reporting
* @returns Status of each server that was started
*/
export async function warmupLspServers(cwd: string, options?: LspWarmupOptions): Promise<LspWarmupResult> {
const config = loadConfig(cwd);
setIdleTimeout(config.idleTimeoutMs);
const servers: LspWarmupResult["servers"] = [];
const lspServers = getLspServers(config);
// Notify caller which servers we're connecting to
if (lspServers.length > 0 && options?.onConnecting) {
options.onConnecting(lspServers.map(([name]) => name));
}
// Start all detected servers in parallel with a short timeout
// Servers that don't respond quickly will be initialized lazily on first use
const results = await Promise.allSettled(
lspServers.map(async ([name, serverConfig]) => {
const client = await getOrCreateClient(serverConfig, cwd, serverConfig.warmupTimeoutMs ?? WARMUP_TIMEOUT_MS);
return { name, client, fileTypes: serverConfig.fileTypes };
}),
);
for (let i = 0; i < results.length; i++) {
const result = results[i];
const [name, serverConfig] = lspServers[i];
if (result.status === "fulfilled") {
servers.push({
name: result.value.name,
status: "ready",
fileTypes: result.value.fileTypes,
});
} else {
const errorMsg = result.reason?.message ?? String(result.reason);
logger.warn("LSP server failed to start", { server: name, error: errorMsg });
servers.push({
name,
status: "error",
fileTypes: serverConfig.fileTypes,
error: errorMsg,
});
}
}
return { servers };
}
/**
* Get status of currently active LSP servers.
*/
export function getLspStatus(): LspServerStatus[] {
return getActiveClients();
}
/**
* Sync in-memory file content to all applicable LSP servers.
* Sends didOpen (if new) or didChange (if already open).
*
* @param absolutePath - Absolute path to the file
* @param content - The new file content
* @param cwd - Working directory for LSP config resolution
* @param servers - Servers to sync to
*/
export async function syncFileContent(
absolutePath: string,
content: string,
cwd: string,
servers: Array<[string, ServerConfig]>,
signal?: AbortSignal,
createMissing = true,
): Promise<void> {
throwIfAborted(signal);
await Promise.allSettled(
servers.map(async ([_serverName, serverConfig]) => {
throwIfAborted(signal);
if (serverConfig.createClient) {
return;
}
const client = createMissing
? await getOrCreateClient(serverConfig, cwd, undefined, signal)
: await getActiveOrPendingClient(serverConfig, cwd, signal);
if (!client) return;
throwIfAborted(signal);
await syncContent(client, absolutePath, content, signal);
}),
);
throwIfAborted(signal);
}
/**
* Notify all LSP servers that a file was saved.
* Assumes content was already synced via syncFileContent.
*
* @param absolutePath - Absolute path to the file
* @param cwd - Working directory for LSP config resolution
* @param servers - Servers to notify
*/
export async function notifyFileSaved(
absolutePath: string,
cwd: string,
servers: Array<[string, ServerConfig]>,
signal?: AbortSignal,
createMissing = true,
): Promise<void> {
throwIfAborted(signal);
await Promise.allSettled(
servers.map(async ([_serverName, serverConfig]) => {
throwIfAborted(signal);
if (serverConfig.createClient) {
return;
}
const client = createMissing
? await getOrCreateClient(serverConfig, cwd, undefined, signal)
: await getActiveOrPendingClient(serverConfig, cwd, signal);
if (!client) return;
await notifySaved(client, absolutePath, signal);
}),
);
throwIfAborted(signal);
}
// Cache config per cwd to avoid repeated file I/O
export const configCache = new Map<string, LspConfig>();
export function getConfig(cwd: string): LspConfig {
let config = configCache.get(cwd);
if (!config) {
config = loadConfig(cwd);
setIdleTimeout(config.idleTimeoutMs);
configCache.set(cwd, config);
}
return config;
}
function isCustomLinter(serverConfig: ServerConfig): boolean {
return Boolean(serverConfig.createClient);
}
export function splitServers(servers: Array<[string, ServerConfig]>): {
lspServers: Array<[string, ServerConfig]>;
customLinterServers: Array<[string, ServerConfig]>;
} {
const lspServers: Array<[string, ServerConfig]> = [];
const customLinterServers: Array<[string, ServerConfig]> = [];
for (const entry of servers) {
if (isCustomLinter(entry[1])) {
customLinterServers.push(entry);
} else {
lspServers.push(entry);
}
}
return { lspServers, customLinterServers };
}
export function getLspServers(config: LspConfig): Array<[string, ServerConfig]> {
return (Object.entries(config.servers) as Array<[string, ServerConfig]>).filter(
([, serverConfig]) => !isCustomLinter(serverConfig),
);
}
export function getLspServersForFile(config: LspConfig, filePath: string): Array<[string, ServerConfig]> {
return getServersForFile(config, filePath).filter(([, serverConfig]) => !isCustomLinter(serverConfig));
}
export function getLspServerForFile(config: LspConfig, filePath: string): [string, ServerConfig] | null {
const servers = getLspServersForFile(config, filePath);
return servers.length > 0 ? servers[0] : null;
}
export function isProjectAwareLspServer(serverConfig: ServerConfig): boolean {
return !serverConfig.createClient && !serverConfig.isLinter;
}
/** True when an LSP error indicates the server doesn't implement the requested method. */
export function isMethodNotFoundError(err: unknown): boolean {
if (!(err instanceof Error)) return false;
const msg = err.message.toLowerCase();
return (
msg.includes("method not found") ||
msg.includes("unhandled method") ||
msg.includes("not supported") ||
msg.includes("-32601")
);
}
export async function reloadServer(client: LspClient, serverName: string, signal?: AbortSignal): Promise<string> {
throwIfAborted(signal);
// rust-analyzer exposes a real reload request. Every other server rejects it
// with method-not-found — that alone justifies the generic fallback. A caller
// cancel or tool timeout must propagate, never be mistaken for an unsupported
// method and swallowed into a bogus "Restarted" (issue #6369).
try {
await sendRequest(client, "rust-analyzer/reloadWorkspace", null, signal);
return `Reloaded ${serverName}`;
} catch (err) {
throwIfAborted(signal);
if (!isMethodNotFoundError(err)) throw err;
// Method not supported — fall through to the generic reload.
}
// 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: {} }, signal);
return `Reloaded ${serverName}`;
} catch {
throwIfAborted(signal);
// The reload notification could not be delivered — the connection is
// wedged or the process already died. Tear the client down (removing it
// from the registry by identity and awaiting confirmed process exit) so
// the next request cold-starts a fresh client. A kill that never confirms
// exit is not a restart: surface the teardown failure truthfully.
//
// On a broker-shared link a per-session teardown only detaches this
// process while the wedged server keeps serving everyone else — ask the
// mux to kill the shared server first (best-effort; it also severs us).
if (client.proc.sharedMux) {
await sendNotification(client, MUX_RESTART_METHOD, undefined, AbortSignal.timeout(2_000)).catch(() => {});
}
if (!(await shutdownClientInstance(client))) {
throw new Error(`Failed to restart ${serverName}: server process did not exit after kill`);
}
return `Restarted ${serverName}`;
}
}
File diff suppressed because it is too large Load Diff
@@ -0,0 +1,170 @@
import * as fs from "node:fs";
import path from "node:path";
import { ToolAbortError, throwIfAborted } from "../tools/tool-errors";
/** Project type detection result */
interface ProjectType {
type: "rust" | "typescript" | "go" | "python" | "unknown";
command?: string[];
description: string;
}
/** Convert a `go.work` use directory into the package pattern `go build` needs. */
function goWorkspaceBuildPattern(diskPath: string): string | null {
const trimmed = diskPath.trim();
if (!trimmed) return null;
const isAbsolute = path.isAbsolute(trimmed) || path.win32.isAbsolute(trimmed);
const normalized = trimmed.replaceAll("\\", "/").replace(/\/+$/, "");
const dir = normalized || ".";
if (dir === ".") return "./...";
if (dir.endsWith("/...")) return dir;
if (isAbsolute || dir.startsWith("./") || dir.startsWith("../")) return `${dir}/...`;
return `./${dir}/...`;
}
/** Parse `go work edit -json` output into per-module package patterns. */
function parseGoWorkspaceBuildPatterns(output: string): string[] {
let parsed: unknown;
try {
parsed = JSON.parse(output);
} catch {
return [];
}
if (!parsed || typeof parsed !== "object" || !("Use" in parsed) || !Array.isArray(parsed.Use)) return [];
const patterns = new Set<string>();
for (const entry of parsed.Use) {
if (!entry || typeof entry !== "object" || !("DiskPath" in entry) || typeof entry.DiskPath !== "string") {
continue;
}
const pattern = goWorkspaceBuildPattern(entry.DiskPath);
if (pattern) patterns.add(pattern);
}
return [...patterns];
}
/** Resolve the `go build` command for a `go.work` workspace. */
async function resolveGoWorkspaceDiagnosticsCommand(cwd: string, signal?: AbortSignal): Promise<string[]> {
const fallback = ["go", "build", "./..."];
try {
const proc = Bun.spawn(["go", "work", "edit", "-json"], {
cwd,
stdout: "pipe",
stderr: "pipe",
windowsHide: true,
});
const abortHandler = () => {
proc.kill();
};
if (signal) {
signal.addEventListener("abort", abortHandler, { once: true });
}
try {
const [stdout] = await Promise.all([new Response(proc.stdout).text(), new Response(proc.stderr).text()]);
const exitCode = await proc.exited;
throwIfAborted(signal);
if (exitCode !== 0) return fallback;
const patterns = parseGoWorkspaceBuildPatterns(stdout);
return patterns.length > 0 ? ["go", "build", ...patterns] : fallback;
} finally {
signal?.removeEventListener("abort", abortHandler);
}
} catch {
if (signal?.aborted) {
throw new ToolAbortError();
}
return fallback;
}
}
/** Detect project type from root markers */
async function detectProjectType(cwd: string, signal?: AbortSignal): Promise<ProjectType> {
// Check for Rust (Cargo.toml)
if (fs.existsSync(path.join(cwd, "Cargo.toml"))) {
return { type: "rust", command: ["cargo", "check", "--message-format=short"], description: "Rust (cargo check)" };
}
// Check for TypeScript (tsconfig.json)
if (fs.existsSync(path.join(cwd, "tsconfig.json"))) {
return { type: "typescript", command: ["npx", "tsc", "--noEmit"], description: "TypeScript (tsc --noEmit)" };
}
// Check for Go workspaces before single-module Go projects.
if (fs.existsSync(path.join(cwd, "go.work"))) {
return {
type: "go",
command: await resolveGoWorkspaceDiagnosticsCommand(cwd, signal),
description: "Go workspace (go build)",
};
}
// Check for Go (go.mod)
if (fs.existsSync(path.join(cwd, "go.mod"))) {
return { type: "go", command: ["go", "build", "./..."], description: "Go (go build)" };
}
// Check for Python (pyproject.toml or pyrightconfig.json)
if (fs.existsSync(path.join(cwd, "pyproject.toml")) || fs.existsSync(path.join(cwd, "pyrightconfig.json"))) {
return { type: "python", command: ["pyright"], description: "Python (pyright)" };
}
return { type: "unknown", description: "Unknown project type" };
}
/** Run workspace diagnostics command and parse output */
export async function runWorkspaceDiagnostics(
cwd: string,
signal?: AbortSignal,
): Promise<{ output: string; projectType: ProjectType }> {
throwIfAborted(signal);
const projectType = await detectProjectType(cwd, signal);
if (!projectType.command) {
return {
output: `Cannot detect project type. Supported: Rust (Cargo.toml), TypeScript (tsconfig.json), Go (go.work/go.mod), Python (pyproject.toml)`,
projectType,
};
}
try {
const proc = Bun.spawn(projectType.command, {
cwd,
stdout: "pipe",
stderr: "pipe",
windowsHide: true,
});
const abortHandler = () => {
proc.kill();
};
if (signal) {
signal.addEventListener("abort", abortHandler, { once: true });
}
try {
const [stdout, stderr] = await Promise.all([
new Response(proc.stdout).text(),
new Response(proc.stderr).text(),
]);
await proc.exited;
throwIfAborted(signal);
const combined = (stdout + stderr).trim();
if (!combined) {
return { output: "No issues found", projectType };
}
// Limit output length
const lines = combined.split("\n");
if (lines.length > 50) {
return { output: `${lines.slice(0, 50).join("\n")}\n[…${lines.length - 50}ln elided…]`, projectType };
}
return { output: combined, projectType };
} finally {
signal?.removeEventListener("abort", abortHandler);
}
} catch (e) {
if (signal?.aborted) {
throw new ToolAbortError();
}
return { output: `Failed to run ${projectType.command.join(" ")}: ${e}`, projectType };
}
}
@@ -0,0 +1,561 @@
import * as fs from "node:fs";
import { isEnoent, logger, once, untilAborted } from "@oh-my-pi/pi-utils";
import type { BunFile } from "bun";
import { FileChangeType, notifyWorkspaceWatchedFiles } from "./client";
import { getServersForFile } from "./config";
import {
captureDiagnosticVersions,
captureOpenFileVersions,
DEFERRED_DIAGNOSTICS_WAIT_TIMEOUT_MS,
type FileDiagnosticsResult,
FileFormatResult,
formatContent,
getDiagnosticsForFile,
INLINE_DIAGNOSTICS_WAIT_TIMEOUT_MS,
limitDiagnosticMessages,
type ServerVersionMap,
} from "./diagnostics";
import { getConfig, notifyFileSaved, splitServers, syncFileContent } from "./servers";
import type { ServerConfig } from "./types";
import { summarizeDiagnosticMessages } from "./utils";
/** Options for creating the LSP writethrough callback */
export interface WritethroughOptions {
/** Whether to format the file using LSP after writing */
enableFormat?: boolean;
/** Whether to get LSP diagnostics after writing */
enableDiagnostics?: boolean;
/** Called when diagnostics arrive after the main timeout. */
onDeferredDiagnostics?: (diagnostics: FileDiagnosticsResult) => void;
/** Signal to cancel a pending deferred diagnostics fetch. */
deferredSignal?: AbortSignal;
/** Transform diagnostics before surfacing them after a successful fetch. */
transformDiagnostics?: (absPath: string, result: FileDiagnosticsResult) => FileDiagnosticsResult;
}
/** Internal resolved form of {@link WritethroughOptions} that the writethrough machinery operates on. */
type ResolvedWritethroughOptions = {
enableFormat: boolean;
enableDiagnostics: boolean;
transformDiagnostics?: (absPath: string, result: FileDiagnosticsResult) => FileDiagnosticsResult;
};
/** Per-file deferred LSP diagnostics wiring for {@link WritethroughCallback}. */
export type WritethroughDeferredHandle = {
onDeferredDiagnostics: (diagnostics: FileDiagnosticsResult) => void;
signal: AbortSignal;
finalize: (diagnostics: FileDiagnosticsResult | undefined) => void;
};
/** Callback type for the LSP writethrough */
export type WritethroughCallback = (
dst: string,
content: string,
signal?: AbortSignal,
file?: BunFile,
batch?: LspWritethroughBatchRequest,
getDeferred?: (dst: string) => WritethroughDeferredHandle | undefined,
) => Promise<FileDiagnosticsResult | undefined>;
/** No-op writethrough callback */
export async function writethroughNoop(
dst: string,
content: string,
_signal?: AbortSignal,
file?: BunFile,
_batch?: LspWritethroughBatchRequest,
_getDeferred?: (dst: string) => WritethroughDeferredHandle | undefined,
): Promise<FileDiagnosticsResult | undefined> {
if (file) {
await file.write(content);
} else {
await Bun.write(dst, content);
}
return undefined;
}
interface PendingWritethrough {
dst: string;
file?: BunFile;
changeType: FileChangeType;
}
interface RunLspWritethroughOptions {
contentAlreadyWritten?: boolean;
}
interface LspWritethroughBatchRequest {
id: string;
flush: boolean;
}
interface LspWritethroughBatchState {
entries: Map<string, PendingWritethrough>;
options: ResolvedWritethroughOptions;
}
const writethroughBatches = new Map<string, LspWritethroughBatchState>();
function getOrCreateWritethroughBatch(id: string, options: ResolvedWritethroughOptions): LspWritethroughBatchState {
const existing = writethroughBatches.get(id);
if (existing) {
existing.options.enableFormat ||= options.enableFormat;
existing.options.enableDiagnostics ||= options.enableDiagnostics;
existing.options.transformDiagnostics ??= options.transformDiagnostics;
return existing;
}
const batch: LspWritethroughBatchState = {
entries: new Map<string, PendingWritethrough>(),
options: { ...options },
};
writethroughBatches.set(id, batch);
return batch;
}
export async function flushLspWritethroughBatch(
id: string,
cwd: string,
signal?: AbortSignal,
): Promise<FileDiagnosticsResult | undefined> {
const state = writethroughBatches.get(id);
if (!state) {
return undefined;
}
writethroughBatches.delete(id);
return flushWritethroughBatch(Array.from(state.entries.values()), cwd, state.options, signal);
}
function mergeDiagnostics(
results: Array<FileDiagnosticsResult | undefined>,
options: ResolvedWritethroughOptions,
): FileDiagnosticsResult | undefined {
const messages: string[] = [];
const servers = new Set<string>();
let hasResults = false;
let hasFormatter = false;
let formatted = false;
for (const result of results) {
if (!result) continue;
hasResults = true;
if (result.server) {
for (const server of result.server.split(",")) {
const trimmed = server.trim();
if (trimmed) {
servers.add(trimmed);
}
}
}
if (result.messages.length > 0) {
messages.push(...result.messages);
}
if (result.formatter !== undefined) {
hasFormatter = true;
if (result.formatter === FileFormatResult.FORMATTED) {
formatted = true;
}
}
}
if (!hasResults && !hasFormatter) {
return undefined;
}
let summary = options.enableDiagnostics ? "no issues" : "OK";
let errored = false;
let limitedMessages = messages;
if (messages.length > 0) {
const summaryInfo = summarizeDiagnosticMessages(messages);
summary = summaryInfo.summary;
errored = summaryInfo.errored;
limitedMessages = limitDiagnosticMessages(messages);
}
const formatter = hasFormatter ? (formatted ? FileFormatResult.FORMATTED : FileFormatResult.UNCHANGED) : undefined;
return {
server: servers.size > 0 ? Array.from(servers).join(", ") : undefined,
messages: limitedMessages,
summary,
errored,
formatter,
};
}
async function scheduleDeferredDiagnosticsFetch(args: {
dst: string;
cwd: string;
servers: Array<[string, ServerConfig]>;
minVersions: ServerVersionMap | undefined;
expectedDocumentVersions: ServerVersionMap | undefined;
signal: AbortSignal;
callback: (diagnostics: FileDiagnosticsResult) => void;
}): Promise<void> {
try {
const deferredTimeout = AbortSignal.timeout(25_000);
const combined = AbortSignal.any([args.signal, deferredTimeout]);
const diagnostics = await getDiagnosticsForFile(args.dst, args.cwd, args.servers, {
signal: combined,
minVersions: args.minVersions,
expectedDocumentVersions: args.expectedDocumentVersions,
timeoutMs: DEFERRED_DIAGNOSTICS_WAIT_TIMEOUT_MS,
});
if (args.signal.aborted || diagnostics === undefined) return;
args.callback(diagnostics);
} catch {
// Cancelled or LSP gave up; silently discard.
}
}
/**
* Fetch post-write diagnostics without making the edit/write block on a slow
* language server.
*
* Blocks inline only briefly ({@link INLINE_DIAGNOSTICS_WAIT_TIMEOUT_MS}) for a
* fresh result. Freshness is enforced by the pre-edit `minVersions` baseline:
* exact document-version matches return immediately, and unversioned/mismatched
* publishes must settle with no newer publish before inline acceptance. If
* nothing fresh arrives in the inline window and a deferred
* channel is available, the in-flight fetch is handed off to deliver late via
* `onDeferredDiagnostics`, and this returns `undefined` so the tool result
* lands immediately. Without a deferred channel (direct/CI callers) it blocks
* for the standard budget so the result is still returned inline.
*/
async function fetchDiagnosticsWithDeferral(args: {
dst: string;
cwd: string;
servers: Array<[string, ServerConfig]>;
minVersions: ServerVersionMap | undefined;
expectedDocumentVersions: ServerVersionMap | undefined;
transformDiagnostics?: ResolvedWritethroughOptions["transformDiagnostics"];
deferred?: { onDeferredDiagnostics: (diagnostics: FileDiagnosticsResult) => void; signal: AbortSignal };
signal?: AbortSignal;
}): Promise<FileDiagnosticsResult | undefined> {
const { dst, cwd, servers, minVersions, expectedDocumentVersions, transformDiagnostics, deferred, signal } = args;
const apply = (d: FileDiagnosticsResult | undefined) =>
d && transformDiagnostics ? transformDiagnostics(dst, d) : d;
if (!deferred) {
// No late-injection channel: block for the standard budget and return inline.
return apply(
await getDiagnosticsForFile(dst, cwd, servers, {
signal,
minVersions,
expectedDocumentVersions,
}),
);
}
// One background fetch with a generous inner budget; await it only briefly inline.
const fetchPromise = getDiagnosticsForFile(dst, cwd, servers, {
signal: deferred.signal,
minVersions,
expectedDocumentVersions,
timeoutMs: DEFERRED_DIAGNOSTICS_WAIT_TIMEOUT_MS,
});
const INLINE_TIMEOUT = Symbol("inline-diagnostics-timeout");
const raced = await Promise.race([
fetchPromise,
Bun.sleep(INLINE_DIAGNOSTICS_WAIT_TIMEOUT_MS).then(() => INLINE_TIMEOUT),
]);
if (raced !== INLINE_TIMEOUT) {
return apply(raced as FileDiagnosticsResult | undefined);
}
// Slow server: deliver late via the deferred channel; nothing inline. The
// deferred sink (edit tool) applies its own dedup, so pass the raw result.
void fetchPromise
.then(diagnostics => {
if (diagnostics && !deferred.signal.aborted) deferred.onDeferredDiagnostics(diagnostics);
})
.catch(() => {});
return undefined;
}
async function runLspWritethrough(
dst: string,
content: string,
cwd: string,
options: ResolvedWritethroughOptions,
changeType: FileChangeType,
signal?: AbortSignal,
file?: BunFile,
deferred?: {
onDeferredDiagnostics: (diagnostics: FileDiagnosticsResult) => void;
signal: AbortSignal;
},
runOptions?: RunLspWritethroughOptions,
): Promise<FileDiagnosticsResult | undefined> {
const { enableFormat, enableDiagnostics } = options;
const contentAlreadyWritten = runOptions?.contentAlreadyWritten ?? false;
let finalContent = content;
const writeContent = async (value: string) => (file ? file.write(value) : Bun.write(dst, value));
const getWritePromise = once(() =>
contentAlreadyWritten && finalContent === content ? Promise.resolve() : writeContent(finalContent),
);
let writeNotified = false;
const notifyWriteCommitted = async (notifySignal: AbortSignal | undefined = signal) => {
if (writeNotified) return;
writeNotified = true;
try {
await notifyWorkspaceWatchedFiles(cwd, [{ filePath: dst, type: changeType }], notifySignal);
} catch (error) {
if (notifySignal?.aborted && !signal?.aborted) {
// The operation budget died mid-notify while the caller is still
// live: allow the post-write retry below to re-announce with the
// caller's signal (didChangeWatchedFiles is idempotent).
writeNotified = false;
return;
}
throw error;
}
};
if (!enableFormat && !enableDiagnostics) {
await getWritePromise();
await notifyWriteCommitted();
return undefined;
}
const config = getConfig(cwd);
const servers = getServersForFile(config, dst);
if (servers.length === 0) {
await getWritePromise();
await notifyWriteCommitted();
return undefined;
}
const { lspServers, customLinterServers } = splitServers(servers);
const useCustomFormatter = enableFormat && customLinterServers.length > 0;
// Capture diagnostic versions BEFORE syncing to detect stale diagnostics
// Bound client creation by the writethrough budget: a hung/broken server
// must not add its full init wait (30s default) to every edit.
const minVersionsPromise = enableDiagnostics ? captureDiagnosticVersions(cwd, servers, 5_000, signal) : undefined;
let minVersions = useCustomFormatter ? undefined : await minVersionsPromise;
let expectedDocumentVersions: ServerVersionMap | undefined;
let formatter: FileFormatResult | undefined;
let diagnostics: FileDiagnosticsResult | undefined;
let timedOut = false;
let synced = false;
let operationSignal: AbortSignal | undefined;
try {
const timeoutSignal = AbortSignal.timeout(5_000);
timeoutSignal.addEventListener(
"abort",
() => {
timedOut = true;
},
{ once: true },
);
operationSignal = signal ? AbortSignal.any([signal, timeoutSignal]) : timeoutSignal;
await untilAborted(operationSignal, async () => {
if (useCustomFormatter) {
// Custom linters operate on on-disk input; the shared pre-write also
// supports implementations that inspect the file before formatting.
if (!contentAlreadyWritten) await writeContent(content);
const [formattedContent, capturedVersions] = await Promise.all([
formatContent(dst, content, cwd, customLinterServers, operationSignal),
minVersionsPromise,
]);
finalContent = formattedContent;
minVersions = capturedVersions;
formatter = finalContent !== content ? FileFormatResult.FORMATTED : FileFormatResult.UNCHANGED;
if (!contentAlreadyWritten || finalContent !== content) await writeContent(finalContent);
await notifyWriteCommitted(operationSignal);
await syncFileContent(dst, finalContent, cwd, lspServers, operationSignal, enableDiagnostics);
} else {
// 1. Sync original content to LSP servers
await syncFileContent(dst, content, cwd, lspServers, operationSignal);
// 2. Format in-memory via LSP
if (enableFormat) {
finalContent = await formatContent(dst, content, cwd, lspServers, operationSignal);
formatter = finalContent !== content ? FileFormatResult.FORMATTED : FileFormatResult.UNCHANGED;
}
// 3. If formatted, sync formatted content to LSP servers
if (finalContent !== content) {
await syncFileContent(dst, finalContent, cwd, lspServers, operationSignal);
}
// 4. Write to disk
await getWritePromise();
await notifyWriteCommitted(operationSignal);
}
if (enableDiagnostics) {
expectedDocumentVersions = await captureOpenFileVersions(dst, cwd, lspServers, operationSignal);
}
// 5. Notify saved to LSP servers
await notifyFileSaved(dst, cwd, lspServers, operationSignal, !useCustomFormatter || enableDiagnostics);
});
synced = true;
} catch {
if (timedOut) {
formatter = undefined;
diagnostics = undefined;
// Schedule background diagnostic fetch if caller wants deferred results
if (deferred && !deferred.signal.aborted && enableDiagnostics) {
void scheduleDeferredDiagnosticsFetch({
dst,
cwd,
servers,
minVersions,
expectedDocumentVersions,
signal: deferred.signal,
callback: deferred.onDeferredDiagnostics,
});
}
}
await getWritePromise();
// The write above committed even though the operation budget elapsed:
// announce it on the caller's signal — the dead `operationSignal` would
// abort the notify before it ever reaches the server.
await notifyWriteCommitted();
}
if (synced && enableDiagnostics) {
diagnostics = await fetchDiagnosticsWithDeferral({
dst,
cwd,
servers,
minVersions,
expectedDocumentVersions,
transformDiagnostics: options.transformDiagnostics,
deferred,
signal,
});
}
if (formatter !== undefined) {
diagnostics ??= {
server: servers.map(([name]) => name).join(", "),
messages: [],
summary: "OK",
errored: false,
};
diagnostics.formatter = formatter;
}
return diagnostics;
}
async function flushWritethroughBatch(
batch: PendingWritethrough[],
cwd: string,
options: ResolvedWritethroughOptions,
signal?: AbortSignal,
getDeferred?: (dst: string) => WritethroughDeferredHandle | undefined,
): Promise<FileDiagnosticsResult | undefined> {
if (batch.length === 0) {
return undefined;
}
const results: Array<FileDiagnosticsResult | undefined> = [];
for (const entry of batch) {
const bundle = getDeferred?.(entry.dst);
let content: string;
try {
content = await fs.promises.readFile(entry.dst, "utf8");
} catch (error) {
if (!isEnoent(error)) throw error;
bundle?.finalize(undefined);
continue;
}
const deferredInner =
bundle &&
({
onDeferredDiagnostics: bundle.onDeferredDiagnostics,
signal: bundle.signal,
} as const);
const diag = await runLspWritethrough(
entry.dst,
content,
cwd,
options,
entry.changeType,
signal,
entry.file,
deferredInner,
{ contentAlreadyWritten: true },
);
bundle?.finalize(diag);
results.push(diag);
}
return mergeDiagnostics(results, options);
}
/** Create a writethrough callback for LSP aware write operations */
export function createLspWritethrough(cwd: string, options?: WritethroughOptions): WritethroughCallback {
const resolvedOptions: ResolvedWritethroughOptions = {
enableFormat: options?.enableFormat ?? false,
enableDiagnostics: options?.enableDiagnostics ?? false,
transformDiagnostics: options?.transformDiagnostics,
};
return async (
dst: string,
content: string,
signal?: AbortSignal,
file?: BunFile,
batch?: LspWritethroughBatchRequest,
getDeferred?: (dst: string) => WritethroughDeferredHandle | undefined,
) => {
const changeType = (await Bun.file(dst).exists()) ? FileChangeType.Changed : FileChangeType.Created;
if (!batch) {
const bundle = getDeferred?.(dst);
const deferredInner =
bundle &&
({
onDeferredDiagnostics: bundle.onDeferredDiagnostics,
signal: bundle.signal,
} as const);
const diagnostics = await runLspWritethrough(
dst,
content,
cwd,
resolvedOptions,
changeType,
signal,
file,
deferredInner,
);
bundle?.finalize(diagnostics);
return diagnostics;
}
// File commits are never deferred: the batch owns only LSP post-processing,
// so a later flush cannot replay an obsolete whole-file snapshot.
try {
await writethroughNoop(dst, content, signal, file);
} catch (error) {
if (batch.flush) {
const pending = writethroughBatches.get(batch.id);
if (pending) {
writethroughBatches.delete(batch.id);
try {
await flushWritethroughBatch(
Array.from(pending.entries.values()),
cwd,
pending.options,
signal,
getDeferred,
);
} catch (flushError) {
logger.warn("Failed to flush pending LSP batch after final write failure", {
batchId: batch.id,
error: flushError instanceof Error ? flushError.message : String(flushError),
});
}
}
}
throw error;
}
const state = getOrCreateWritethroughBatch(batch.id, resolvedOptions);
state.entries.set(dst, { dst, file, changeType });
if (!batch.flush) return undefined;
writethroughBatches.delete(batch.id);
return flushWritethroughBatch(Array.from(state.entries.values()), cwd, state.options, signal, getDeferred);
};
}