返回源码地图

packages/lsp/lsp-stdio/src/connection.ts

main snapshot · da00f7f5358f · 正文引用章节 10;完整原文可核对,不声称全文件人工逐行审计

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

1/**
2 * A JSON-RPC endpoint over one language server spawned through the subprocess
3 * capability. Owns id correlation, outbound requests/notifications, and inbound
4 * server→client requests: it answers `workspace/configuration` from static
5 * config, and rejects `workspace/applyEdit` (this host never applies edits or
6 * runs commands). It caps stderr, surfaces framing/decoder failures as a
7 * fatal close, and exposes managed-range termination through the handle so the
8 * instance owns teardown; platform mechanics live in the subprocess
9 * Service Provider.
10 * @module @deepseek-ai/dsh-lsp-stdio/connection
11 */
12
13import type { Writable } from 'node:stream'
14import type { SubprocessHandle, SubprocessSpawnSpec } from '@deepseek-ai/dsh-subprocess'
15import { encodeMessage, MessageDecoder } from './framing.ts'
16
17/** How to launch the server and answer its config requests. */
18export interface ConnectionSpec {
19 /** The resolved absolute executable path (no shell). */
20 readonly command: string
21 /** Arguments passed to the executable. */
22 readonly args: readonly string[]
23 /** The child's working directory (the canonical workspace). */
24 readonly cwd: string
25 /** Explicit child environment overrides; the subprocess provider owns its ambient scrub. */
26 readonly env: Record<string, string>
27 /** Largest single framed message accepted from the server. */
28 readonly maxMessageBytes: number
29 /** Largest stderr tail retained for diagnostics. */
30 readonly maxStderrBytes: number
31 /**
32 * The subprocess spec's `graceMs`: the SIGTERM→SIGKILL window of
33 * {@link LspConnection.terminate}'s escalation, and the bound for draining
34 * pipes a surviving helper still holds after the server exits.
35 */
36 readonly killGraceMs: number
37 /** Static answer to every `workspace/configuration` item. */
38 readonly configuration: unknown
39}
40
41interface Pending {
42 resolve: (value: unknown) => void
43 reject: (error: Error) => void
44}
45
46/**
47 * Write one JSON-RPC message to the child stdin.
48 * @param stdin - the spawned server stdin.
49 * @param message - the unencoded JSON-RPC message.
50 * @param done - callback that reports asynchronous stream settlement.
51 */
52export type ConnectionWriter = (
53 stdin: Writable,
54 message: unknown,
55 done: (error?: Error | null) => void,
56) => void
57
58/** Spawn one subprocess for this connection (the provider passes `ctx.subprocess.spawn`). */
59export type ConnectionSpawner = (spec: SubprocessSpawnSpec) => SubprocessHandle
60
61const writeConnectionMessage: ConnectionWriter = (stdin, message, done) => {
62 stdin.write(encodeMessage(message), done)
63}
64
65/** A live JSON-RPC endpoint bound to one child process. */
66export class LspConnection {
67 private readonly handle: SubprocessHandle
68 private readonly stdin: Writable
69 private readonly decoder: MessageDecoder
70 private readonly pending = new Map<number, Pending>()
71 private nextId = 1
72 private closeReason: Error | undefined
73 /** Set once the process has fully exited; the instance awaits it during teardown. */
74 readonly closed: Promise<void>
75
76 /**
77 * @param spec - how to launch the server and answer its config requests.
78 * @param spawner - the subprocess seam's spawn (the provider passes `ctx.subprocess.spawn`).
79 * @param onServerRequest - answers a server→client request; rejects to send an error response.
80 * @param writer - message writer; tests inject callback failures without relying on OS pipe races.
81 */
82 constructor(
83 spec: ConnectionSpec,
84 spawner: ConnectionSpawner,
85 private readonly onServerRequest: (method: string, params: unknown) => Promise<unknown>,
86 private readonly writer: ConnectionWriter = writeConnectionMessage,
87 ) {
88 this.decoder = new MessageDecoder(spec.maxMessageBytes)
89 // stdin/stdout are piped protocol streams this endpoint frames itself;
90 // stderr is a collected diagnostic tail (no spill — the bounded tail IS
91 // the contract). The seam owns managed-range signalling and observation.
92 this.handle = spawner({
93 argv: [spec.command, ...spec.args],
94 cwd: spec.cwd,
95 stdio: {
96 stdin: 'pipe',
97 stdout: 'pipe',
98 stderr: { maxBytes: spec.maxStderrBytes },
99 },
100 graceMs: spec.killGraceMs,
101 // The seam merges explicit config entries after its ambient scrub, so a
102 // configured credential or DSH_* fact reaches the child deliberately.
103 env: spec.env,
104 })
105 /* v8 ignore start -- 'pipe' dispositions expose both streams by the seam contract; defensive. */
106 if (this.handle.stdin === undefined || this.handle.stdout === undefined) {
107 throw new Error('lsp-stdio: subprocess implementation dropped a piped protocol stream')
108 }
109 /* v8 ignore stop */
110 this.stdin = this.handle.stdin
111 this.closed = new Promise<void>((resolve) => {
112 const close = (): void => {
113 const reason = this.closeReason ?? new Error(this.exitMessage())
114 // Record the reason so any request issued AFTER close rejects immediately instead of hanging
115 // (a closed process sends no further responses).
116 this.closeReason = reason
117 this.failAll(reason)
118 resolve()
119 }
120 this.handle.done.then(close, (error: unknown) => {
121 // A spawn-level failure never produces a close event; the rejection is
122 // the fatal cause and the close boundary at once.
123 this.fail(asError(error))
124 close()
125 })
126 })
127 // Child stdin can fail while the process itself remains alive (for example, a server closes fd
128 // 0). Treat that as a fatal connection error so pending requests reject immediately instead of
129 // waiting for a process-close event that may never arrive.
130 this.stdin.on('error', (error) => { this.fail(error) })
131 this.handle.stdout.on('data', (chunk: Buffer) => { this.onStdout(chunk) })
132 }
133
134 /** The retained stderr tail, for diagnostics on a failed server. */
135 get stderrTail(): string {
136 /* v8 ignore next -- the collect disposition always exposes a stderr reader; defensive. */
137 return this.handle.collected.stderr?.readFrom(0).text ?? ''
138 }
139
140 /** Whether the transport has failed even if the child close event has not arrived yet. */
141 get failed(): boolean {
142 return this.closeReason !== undefined
143 }
144
145 /**
146 * Test whether a caught error is this connection's retained fatal transport cause.
147 * @param error - error caught by the instance or provider.
148 * @returns `true` only when this connection produced that exact failure.
149 */
150 failedWith(error: unknown): boolean {
151 return this.closeReason === error
152 }
153
154 /**
155 * Send a request and await its result.
156 * @param method - the JSON-RPC method.
157 * @param params - the request params.
158 * @returns the response result; rejects on an error response, write failure, or close.
159 */
160 request(method: string, params: unknown): Promise<unknown> {
161 const id = this.nextId++
162 const promise = new Promise<unknown>((resolve, reject) => {
163 if (this.closeReason !== undefined) {
164 reject(this.closeReason)
165 return
166 }
167 this.pending.set(id, { resolve, reject })
168 // `write()` records either synchronous or callback-delivered failures on the connection and
169 // rejects every pending request. This handler only consumes the write promise itself.
170 void this.write({ jsonrpc: '2.0', id, method, params }).catch(() => {})
171 })
172 // A caller that stops awaiting (e.g. an aborted query) can leave this promise to reject later
173 // when the process closes; a benign no-op handler keeps that from surfacing as an unhandled
174 // rejection. The returned promise still delivers the rejection to the caller's own await/catch.
175 promise.catch(() => {})
176 return promise
177 }
178
179 /**
180 * Send a notification (no id, no response).
181 * @param method - the JSON-RPC method.
182 * @param params - the notification params.
183 * @returns a promise that settles when the framed notification has been written.
184 */
185 notify(method: string, params: unknown): Promise<void> {
186 return this.write({ jsonrpc: '2.0', method, params })
187 }
188
189 /**
190 * Send a `$/cancelRequest` for an in-flight request id (best-effort; ignores write failure).
191 * @param requestId - the numeric id of the request to cancel.
192 */
193 cancel(requestId: number): void {
194 // The server is already gone or unwritable when this rejects; `write()` has recorded the fatal
195 // connection failure and rejected the pending request, so cancellation remains best-effort.
196 void this.write({ jsonrpc: '2.0', method: '$/cancelRequest', params: { id: requestId } }).catch(() => {})
197 }
198
199 /**
200 * The id the NEXT `request()` will use, so the instance can pre-arm a cancel.
201 * @returns the numeric id the next request will be assigned.
202 */
203 peekNextId(): number {
204 return this.nextId
205 }
206
207 /** Terminate the server's provider-managed range (idempotent). */
208 terminate(): void {
209 this.handle.terminate()
210 }
211
212 /**
213 * Wait until the owned managed range is empty.
214 * @param signal - optional bound for the wait.
215 * @returns `true` when the range is empty, or `false` when the signal aborted first.
216 */
217 async waitForManagedRangeExit(signal?: AbortSignal): Promise<boolean> {
218 return await this.handle.waitForExit(signal)
219 }
220
221 private onStdout(chunk: Buffer): void {
222 let messages: unknown[]
223 try {
224 messages = this.decoder.push(chunk)
225 } catch (error) {
226 // A framing/JSON failure corrupts the stream position irrecoverably: fail the instance and
227 // terminate the managed range so helper processes do not outlive the leader.
228 this.fail(asError(error))
229 this.handle.terminate()
230 return
231 }
232 for (const message of messages) this.dispatch(message)
233 }
234
235 private dispatch(message: unknown): void {
236 if (message === null || typeof message !== 'object') return
237 const frame = message as Record<string, unknown>
238 const id = frame.id
239 const method = frame.method
240 if (typeof method === 'string' && (typeof id === 'number' || typeof id === 'string')) {
241 // A response-write failure has already invalidated the connection in `write()`.
242 /* v8 ignore next -- protocol tests exercise response writes; only a simultaneous connection
243 failure makes this consumption handler run. */
244 void this.handleServerRequest(id, method, frame.params).catch(() => {})
245 return
246 }
247 if (typeof method === 'string') {
248 // A server→client notification (e.g. diagnostics, logs): ignored by this MVP host.
249 return
250 }
251 if (typeof id === 'number') this.handleResponse(id, frame)
252 }
253
254 private async handleServerRequest(id: number | string, method: string, params: unknown): Promise<void> {
255 try {
256 const result = await this.onServerRequest(method, params)
257 await this.write({ jsonrpc: '2.0', id, result })
258 } catch (error) {
259 await this.write({ jsonrpc: '2.0', id, error: { code: -32601, message: asError(error).message } })
260 }
261 }
262
263 private handleResponse(id: number, frame: Record<string, unknown>): void {
264 const pending = this.pending.get(id)
265 if (!pending) return
266 this.pending.delete(id)
267 const error = frame.error
268 if (error !== null && typeof error === 'object') {
269 const record = error as Record<string, unknown>
270 pending.reject(new Error(typeof record.message === 'string' ? record.message : 'LSP error response'))
271 return
272 }
273 pending.resolve(frame.result)
274 }
275
276 private write(message: unknown): Promise<void> {
277 if (this.closeReason !== undefined) return Promise.reject(this.closeReason)
278 return new Promise<void>((resolve, reject) => {
279 const done = (error?: Error | null): void => {
280 if (error === undefined || error === null) {
281 resolve()
282 return
283 }
284 this.fail(error)
285 reject(error)
286 }
287 try {
288 this.writer(this.stdin, message, done)
289 /* v8 ignore start -- Node stream write failures are callback-delivered; this guards a
290 nonconforming Writable implementation throwing synchronously. */
291 } catch (error) {
292 const failure = asError(error)
293 this.fail(failure)
294 reject(failure)
295 }
296 /* v8 ignore stop */
297 })
298 }
299
300 /** The exit-close error message, appending the retained stderr tail when the server wrote any. */
301 private exitMessage(): string {
302 const tail = this.stderrTail.trim()
303 return tail === '' ? 'language server exited' : `language server exited; stderr: ${tail}`
304 }
305
306 private fail(error: Error): void {
307 /* v8 ignore next -- the second arm (closeReason already set) needs two fail() calls before close; defensive. */
308 if (this.closeReason === undefined) this.closeReason = error
309 this.failAll(error)
310 }
311
312 private failAll(error: Error): void {
313 const waiting = [...this.pending.values()]
314 this.pending.clear()
315 for (const pending of waiting) pending.reject(error)
316 }
317}
318
319/** Coerce an unknown thrown value to an `Error`. */
320function asError(value: unknown): Error {
321 /* v8 ignore next -- the non-Error branch guards against a non-Error throw, which our paths never produce. */
322 return value instanceof Error ? value : new Error(String(value))
323}