1
/**2
* A JSON-RPC endpoint over one language server spawned through the subprocess3
* capability. Owns id correlation, outbound requests/notifications, and inbound4
* server→client requests: it answers `workspace/configuration` from static5
* config, and rejects `workspace/applyEdit` (this host never applies edits or6
* runs commands). It caps stderr, surfaces framing/decoder failures as a7
* fatal close, and exposes managed-range termination through the handle so the8
* instance owns teardown; platform mechanics live in the subprocess9
* Service Provider.10
* @module @deepseek-ai/dsh-lsp-stdio/connection11
*/13
import type { Writable } from 'node:stream'14
import type { SubprocessHandle, SubprocessSpawnSpec } from '@deepseek-ai/dsh-subprocess'15
import { encodeMessage, MessageDecoder } from './framing.ts'17
/** How to launch the server and answer its config requests. */18
export interface ConnectionSpec {19
/** The resolved absolute executable path (no shell). */20
readonly command: string21
/** Arguments passed to the executable. */22
readonly args: readonly string[]23
/** The child's working directory (the canonical workspace). */24
readonly cwd: string25
/** 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: number29
/** Largest stderr tail retained for diagnostics. */30
readonly maxStderrBytes: number31
/**32
* The subprocess spec's `graceMs`: the SIGTERM→SIGKILL window of33
* {@link LspConnection.terminate}'s escalation, and the bound for draining34
* pipes a surviving helper still holds after the server exits.35
*/36
readonly killGraceMs: number37
/** Static answer to every `workspace/configuration` item. */38
readonly configuration: unknown39
}41
interface Pending {42
resolve: (value: unknown) => void43
reject: (error: Error) => void44
}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
*/52
export type ConnectionWriter = (53
stdin: Writable,54
message: unknown,55
done: (error?: Error | null) => void,56
) => void58
/** Spawn one subprocess for this connection (the provider passes `ctx.subprocess.spawn`). */59
export type ConnectionSpawner = (spec: SubprocessSpawnSpec) => SubprocessHandle61
const writeConnectionMessage: ConnectionWriter = (stdin, message, done) => {62
stdin.write(encodeMessage(message), done)63
}65
/** A live JSON-RPC endpoint bound to one child process. */66
export class LspConnection {67
private readonly handle: SubprocessHandle68
private readonly stdin: Writable69
private readonly decoder: MessageDecoder70
private readonly pending = new Map<number, Pending>()71
private nextId = 172
private closeReason: Error | undefined73
/** Set once the process has fully exited; the instance awaits it during teardown. */74
readonly closed: Promise<void>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 IS91
// 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 a102
// 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.stdin111
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 hanging115
// (a closed process sends no further responses).116
this.closeReason = reason117
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 is122
// 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 fd128
// 0). Treat that as a fatal connection error so pending requests reject immediately instead of129
// 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
}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
}140
/** Whether the transport has failed even if the child close event has not arrived yet. */141
get failed(): boolean {142
return this.closeReason !== undefined143
}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 === error152
}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
return166
}167
this.pending.set(id, { resolve, reject })168
// `write()` records either synchronous or callback-delivered failures on the connection and169
// 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 later173
// when the process closes; a benign no-op handler keeps that from surfacing as an unhandled174
// rejection. The returned promise still delivers the rejection to the caller's own await/catch.175
promise.catch(() => {})176
return promise177
}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
}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 fatal195
// 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
}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.nextId205
}207
/** Terminate the server's provider-managed range (idempotent). */208
terminate(): void {209
this.handle.terminate()210
}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
}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 and227
// terminate the managed range so helper processes do not outlive the leader.228
this.fail(asError(error))229
this.handle.terminate()230
return231
}232
for (const message of messages) this.dispatch(message)233
}235
private dispatch(message: unknown): void {236
if (message === null || typeof message !== 'object') return237
const frame = message as Record<string, unknown>238
const id = frame.id239
const method = frame.method240
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 connection243
failure makes this consumption handler run. */244
void this.handleServerRequest(id, method, frame.params).catch(() => {})245
return246
}247
if (typeof method === 'string') {248
// A server→client notification (e.g. diagnostics, logs): ignored by this MVP host.249
return250
}251
if (typeof id === 'number') this.handleResponse(id, frame)252
}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
}263
private handleResponse(id: number, frame: Record<string, unknown>): void {264
const pending = this.pending.get(id)265
if (!pending) return266
this.pending.delete(id)267
const error = frame.error268
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
return272
}273
pending.resolve(frame.result)274
}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
return283
}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 a290
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
}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
}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 = error309
this.failAll(error)310
}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
}319
/** Coerce an unknown thrown value to an `Error`. */320
function 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
}