返回源码地图

packages/sdk/client/src/client.ts

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

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

1/**
2 * Low-level JSON-RPC client for a DeepSeek Harness SDK runtime subprocess.
3 * {@link HarnessClient} owns the child process: it spawns the runtime, speaks
4 * the `@deepseek-ai/dsh-sdk-protocol` wire over the child's stdio, fans
5 * server notifications out to subscriptions, and tears the child down to
6 * quiescence through a private EOF → SIGTERM → SIGKILL ladder. The design
7 * twin is the Python SDK's `HarnessClient` (`python/sdk`); both drive the
8 * same runtime protocol. This client runs OUTSIDE any harness context, so it
9 * spawns directly rather than through the `dsh-subprocess` service — the
10 * seam's documented exception for SDK-managed transports.
11 *
12 * @module @deepseek-ai/dsh-sdk-client/client
13 */
14
15import { spawn, type ChildProcess } from 'node:child_process'
16import {
17 JsonRpcLineTransport,
18 JsonRpcResponseError,
19 type InitializeParams,
20 type InitializeResult,
21 type SessionPromptParams,
22 type SdkPromptContentBlock,
23} from '@deepseek-ai/dsh-sdk-protocol'
24import { disposeRuntimeProcess } from './dispose.ts'
25import { resolveDshLaunch, type RuntimeProcessOptions } from './launch.ts'
26import type { HarnessClientOptions, HarnessNotification, NotificationFilter } from './types.ts'
27
28/** Retained stderr lines used to diagnose an unexpected runtime death. */
29const STDERR_TAIL_LIMIT = 400
30
31/** Grace for the runtime's stdio streams to settle after its exit edge. */
32const STREAM_SETTLE_MS = 100
33
34/**
35 * The runtime subprocess is gone or unusable: it exited, its stdio closed, or
36 * it was never launchable. The message carries the exit code and a stderr
37 * tail when available.
38 */
39export class TransportClosedError extends Error {
40 /** @param message - the failure description, including any stderr tail. */
41 constructor(message: string) {
42 super(message)
43 this.name = 'TransportClosedError'
44 }
45}
46
47/** A request exceeded {@link HarnessClientOptions.requestTimeoutMs}. */
48export class RequestTimeoutError extends Error {
49 /** @param message - which method timed out. */
50 constructor(message: string) {
51 super(message)
52 this.name = 'RequestTimeoutError'
53 }
54}
55
56/**
57 * The runtime answered outside its documented protocol (for example a
58 * `session/prompt` response without `accepted: true`).
59 */
60export class SdkProtocolError extends Error {
61 /** @param message - the protocol violation description. */
62 constructor(message: string) {
63 super(message)
64 this.name = 'SdkProtocolError'
65 }
66}
67
68interface SubscriptionState {
69 readonly queue: HarnessNotification[]
70 readonly waiters: { resolve: (item: HarnessNotification) => void; reject: (error: Error) => void }[]
71 readonly filter: NotificationFilter | undefined
72 failure: Error | undefined
73}
74
75/** One client-side notification stream returned by {@link HarnessClient.subscribe}. */
76export interface NotificationSubscription extends AsyncIterable<HarnessNotification> {
77 /**
78 * Await the next matching notification.
79 * @returns the notification; after the runtime died, drains what was
80 * already delivered and then rejects; after {@link close}, rejects
81 * immediately (the queue is dropped).
82 */
83 next(): Promise<HarnessNotification>
84
85 /**
86 * Drain one already-delivered notification without waiting.
87 * @returns the next queued notification, or `undefined` when none is queued.
88 */
89 tryNext(): HarnessNotification | undefined
90
91 /** Detach from the client; queued items drop and pending waiters reject. */
92 close(): void
93}
94
95/** Internal producer side of a public notification subscription. */
96class NotificationSubscriptionImpl implements NotificationSubscription {
97 constructor(
98 private readonly state: SubscriptionState,
99 private readonly unsubscribe: () => void,
100 ) {}
101
102 /**
103 * Await the next matching notification.
104 * @returns the notification; after the runtime died, drains what was
105 * already delivered and then rejects; after {@link close}, rejects
106 * immediately (the queue is dropped).
107 */
108 next(): Promise<HarnessNotification> {
109 const queued = this.state.queue.shift()
110 if (queued !== undefined) return Promise.resolve(queued)
111 if (this.state.failure !== undefined) return Promise.reject(this.state.failure)
112 return new Promise((resolve, reject) => {
113 this.state.waiters.push({ resolve, reject })
114 })
115 }
116
117 /**
118 * Drain one already-delivered notification without waiting.
119 * @returns the next queued notification, or `undefined` when none is queued.
120 */
121 tryNext(): HarnessNotification | undefined {
122 return this.state.queue.shift()
123 }
124
125 /** Detach from the client; queued items drop and pending waiters reject. */
126 close(): void {
127 this.unsubscribe()
128 // The drop is part of this method's contract; a runtime-death fail() keeps
129 // the queue so already-delivered notifications remain drainable.
130 this.state.queue.length = 0
131 this.fail(new TransportClosedError('notification subscription closed'))
132 }
133
134 /**
135 * Reject pending and future waits (delivery stops; the first failure wins).
136 * Already-queued notifications remain drainable via {@link next}/{@link tryNext}.
137 * @param error - the terminal failure delivered to waiters.
138 */
139 fail(error: Error): void {
140 this.state.failure ??= error
141 for (const waiter of this.state.waiters.splice(0)) waiter.reject(this.state.failure)
142 }
143
144 /**
145 * Deliver one notification to a waiter or the queue when the filter
146 * matches. A throwing filter fails only THIS subscription (detached, the
147 * throw becomes its terminal error) — it never disturbs sibling
148 * subscriptions or the transport's read loop.
149 * @param notification - the wire notification to deliver.
150 */
151 push(notification: HarnessNotification): void {
152 let matches: boolean
153 try {
154 matches = this.state.filter === undefined || this.state.filter(notification)
155 } catch (error) {
156 this.unsubscribe()
157 this.fail(error instanceof Error ? error : new Error(String(error)))
158 return
159 }
160 if (!matches) return
161 const waiter = this.state.waiters.shift()
162 if (waiter !== undefined) waiter.resolve(notification)
163 else this.state.queue.push(notification)
164 }
165
166 /**
167 * Iterate notifications until the subscription or runtime closes (the
168 * terminating rejection propagates).
169 * @returns an async iterator over {@link next} results.
170 */
171 async * [Symbol.asyncIterator](): AsyncIterator<HarnessNotification> {
172 for (;;) yield await this.next()
173 }
174}
175
176/**
177 * JSON-RPC client for the DeepSeek Harness SDK runtime over subprocess stdio.
178 *
179 * The subprocess starts lazily on {@link start} and is owned by this instance
180 * until {@link close}, which requests protocol `shutdown` and then walks the
181 * shared EOF → SIGTERM → SIGKILL dispose ladder to quiescence. There is no
182 * wire-level cancel: a timed-out request stays running server-side until the
183 * runtime is closed.
184 */
185export class HarnessClient {
186 /** Original public dsh launch and timeout options for this client. */
187 readonly options: HarnessClientOptions
188 private readonly runtime: RuntimeProcessOptions
189 private child: ChildProcess | undefined
190 private transport: JsonRpcLineTransport | undefined
191 private readonly stderrTail: string[] = []
192 private readonly subscriptions = new Map<string, NotificationSubscriptionImpl>()
193 private readonly sessionParents = new Map<string, string>()
194 private subscriptionSerial = 0
195 private exitCode: number | null | undefined
196 private spawnError: Error | undefined
197 private streamsSettled: Promise<void> = Promise.resolve()
198 private closeTask: Promise<void> | undefined
199
200 /** @param options - dsh profile, patch, home, process, environment, and timeout options. */
201 constructor(options?: HarnessClientOptions)
202 constructor(options: HarnessClientOptions = {}, runtime?: RuntimeProcessOptions) {
203 this.options = options
204 this.runtime = runtime ?? resolveDshLaunch(options)
205 }
206
207 /**
208 * Spawn the runtime subprocess and start reading frames. Idempotent while
209 * the process is live; rejects reuse after {@link close}.
210 */
211 start(): void {
212 if (this.closeTask !== undefined) throw new TransportClosedError('DeepSeek Harness runtime client is closed')
213 if (this.child !== undefined) return
214 const child = spawn(this.runtime.command, this.runtime.args, {
215 cwd: this.runtime.cwd,
216 env: this.runtime.environment(),
217 stdio: ['pipe', 'pipe', 'pipe'],
218 })
219 this.child = child
220 child.once('error', (error) => {
221 this.spawnError = error
222 // A spawn failure destroys the pipes without an input 'end' edge, so the
223 // transport's pending requests must be failed here.
224 this.transport?.close()
225 this.failSubscriptions(this.closedError('DeepSeek Harness runtime failed to start'))
226 })
227 // Writes racing the runtime's death EPIPE on stdin; the exit edge below is
228 // the real signal, so the stream-level error only needs to be non-fatal.
229 // The timing of that race is not deterministically reproducible.
230 /* v8 ignore next */
231 child.stdin.on('error', () => {})
232 let stderrBuffer = ''
233 child.stderr.setEncoding('utf8')
234 child.stderr.on('data', (chunk: string) => {
235 stderrBuffer += chunk
236 const newline = stderrBuffer.lastIndexOf('\n')
237 if (newline >= 0) {
238 this.appendStderr(stderrBuffer.slice(0, newline).split('\n'))
239 stderrBuffer = stderrBuffer.slice(newline + 1)
240 }
241 })
242 let signalStreamsSettled!: () => void
243 this.streamsSettled = new Promise((resolve) => { signalStreamsSettled = resolve })
244 const settled = { stderr: false, exited: false }
245 const maybeSettle = (): void => {
246 if (settled.stderr && settled.exited) signalStreamsSettled()
247 }
248 child.stderr.once('close', () => {
249 if (stderrBuffer.length > 0) this.appendStderr([stderrBuffer])
250 settled.stderr = true
251 maybeSettle()
252 })
253 child.once('exit', (code) => {
254 this.exitCode = code
255 settled.exited = true
256 maybeSettle()
257 this.failSubscriptions(this.closedError('DeepSeek Harness runtime exited'))
258 })
259 child.once('close', () => {
260 // All stdio has settled: stdout 'end' already drained every tail frame,
261 // so closing now cannot drop responses — it only fails requests that
262 // will never be answered.
263 this.transport?.close()
264 })
265 const transport = new JsonRpcLineTransport(child.stdout, child.stdin)
266 transport.onNotification((method, params) => { this.dispatchNotification({ method, params }) })
267 transport.start()
268 this.transport = transport
269 }
270
271 /**
272 * Perform the process-wide handshake.
273 * @param params - workspace cwd plus the provider/model route.
274 * @returns the runtime's wire identity.
275 */
276 async initialize(params: InitializeParams): Promise<InitializeResult> {
277 const result = await this.request('initialize', { ...params }, this.runtime.initializeTimeoutMs)
278 if (!isRecord(result) || !isRecord(result.serverInfo)
279 || typeof result.serverInfo.name !== 'string' || typeof result.serverInfo.version !== 'string') {
280 throw new SdkProtocolError(`initialize returned no server identity: ${JSON.stringify(result)}`)
281 }
282 return { serverInfo: { name: result.serverInfo.name, version: result.serverInfo.version } }
283 }
284
285 /**
286 * Queue one prompt and return its durable inbox identity.
287 * @param sessionId - target session; an unknown id creates it.
288 * @param contentBlocks - the user message, sent verbatim.
289 * @returns the queued message id.
290 */
291 async prompt(sessionId: string, contentBlocks: SdkPromptContentBlock[]): Promise<string> {
292 const params: SessionPromptParams = { sessionId, contentBlocks }
293 const result = await this.request('session/prompt', { ...params })
294 if (!isRecord(result) || typeof result.messageId !== 'string') {
295 throw new SdkProtocolError(`session/prompt returned no message id: ${JSON.stringify(result)}`)
296 }
297 return result.messageId
298 }
299
300 /**
301 * Send one JSON-RPC request and await its result.
302 * @param method - the wire method name.
303 * @param params - the params object; omitted params send `{}`.
304 * @param timeoutMs - per-call override of {@link HarnessClientOptions.requestTimeoutMs}.
305 * @returns the raw result; rejects with {@link JsonRpcResponseError} on a
306 * protocol error response, {@link RequestTimeoutError} on timeout, and
307 * {@link TransportClosedError} when the runtime is gone.
308 */
309 async request(method: string, params?: object, timeoutMs?: number): Promise<unknown> {
310 this.start()
311 // A dead runtime cannot answer; fail with process context instead of
312 // writing into a destroyed pipe and hanging until the timeout.
313 if (this.exitCode !== undefined || this.spawnError !== undefined) {
314 await this.settleStreams()
315 throw this.closedError('DeepSeek Harness runtime is not running')
316 }
317 const transport = this.transport
318 /* v8 ignore next -- start() either sets the transport or throws */
319 if (transport === undefined) throw new TransportClosedError('DeepSeek Harness runtime is not running')
320 const timeout = timeoutMs ?? this.runtime.requestTimeoutMs
321 try {
322 if (timeout === undefined) return await transport.request(method, params ?? {})
323 // The abort signal makes the timeout an abandonment: the transport drops
324 // its pending entry, so repeated bounded requests against a hung method
325 // retain no per-call state (the server-side work still runs to close).
326 const abandon = new AbortController()
327 const timer = setTimeout(() => {
328 const stderr = this.stderrTail.length === 0 ? '' : `; stderr tail:\n${this.stderrTail.join('\n')}`
329 abandon.abort(new RequestTimeoutError(`${method} timed out after ${timeout}ms waiting for ${this.runtime.description}${stderr}`))
330 }, timeout)
331 try {
332 return await transport.request(method, params ?? {}, abandon.signal)
333 } finally {
334 clearTimeout(timer)
335 }
336 } catch (error) {
337 if (error instanceof JsonRpcResponseError || error instanceof RequestTimeoutError) throw error
338 // Transport-level failures gain process context: exit code + stderr tail.
339 await this.settleStreams()
340 throw this.closedError(errorMessage(error))
341 }
342 }
343
344 /**
345 * Subscribe to server notifications.
346 * @param filter - optional predicate; omitted means every notification.
347 * @returns the subscription handle; close it to stop delivery. After
348 * {@link close} or runtime death the handle is born failed — there is no
349 * producer left, so `next()` rejects instead of waiting forever.
350 */
351 subscribe(filter?: NotificationFilter): NotificationSubscription {
352 const id = String(this.subscriptionSerial++)
353 const state: SubscriptionState = { queue: [], waiters: [], filter, failure: undefined }
354 const subscription = new NotificationSubscriptionImpl(state, () => { this.subscriptions.delete(id) })
355 if (this.closeTask !== undefined || this.exitCode !== undefined || this.spawnError !== undefined) {
356 subscription.fail(this.closedError('DeepSeek Harness runtime closed'))
357 return subscription
358 }
359 this.subscriptions.set(id, subscription)
360 return subscription
361 }
362
363 /**
364 * Subscribe to one session and the descendants discovered from
365 * `subagent.started` lineage edges. The runtime notifies for every session
366 * in its context, so this client applies the scope.
367 * @param sessionId - the root session id.
368 * @returns the filtered subscription handle.
369 */
370 subscribeSessionTree(sessionId: string): NotificationSubscription {
371 return this.subscribe((notification) => {
372 const params = notification.params
373 if (notification.method === 'subagent.started' || notification.method === 'subagent.finished') {
374 const parentId = params.parentSessionId
375 if (typeof parentId === 'string' && this.isDescendantOf(parentId, sessionId)) return true
376 return params.childSessionId === sessionId
377 }
378 const relatedId = params.sessionId
379 return typeof relatedId === 'string' && this.isDescendantOf(relatedId, sessionId)
380 })
381 }
382
383 /**
384 * Shut the runtime down and reap it: a best-effort protocol `shutdown`
385 * bounded by `shutdownTimeoutMs`, then the shared stdin-EOF → SIGTERM →
386 * SIGKILL ladder until the process actually exited. Idempotent.
387 * @returns settlement of the complete teardown.
388 */
389 close(): Promise<void> {
390 this.closeTask ??= this.performClose()
391 return this.closeTask
392 }
393
394 private async performClose(): Promise<void> {
395 const child = this.child
396 if (child === undefined) return
397 try {
398 await this.request('shutdown', undefined, this.runtime.shutdownTimeoutMs ?? 1_000)
399 } catch (error) {
400 // Diagnostic only: the dispose ladder below is the authoritative teardown
401 // for a runtime that cannot answer shutdown anymore.
402 this.appendStderr([`shutdown request failed: ${errorMessage(error)}`])
403 }
404 await disposeRuntimeProcess(child, {
405 disposeEofGraceMs: this.runtime.disposeEofGraceMs ?? 6_000,
406 disposeGraceMs: this.runtime.disposeGraceMs ?? 3_000,
407 })
408 this.transport?.close()
409 this.failSubscriptions(this.closedError('DeepSeek Harness runtime closed'))
410 }
411
412 private dispatchNotification(notification: HarnessNotification): void {
413 this.recordSessionRelationship(notification)
414 for (const subscription of this.subscriptions.values()) subscription.push(notification)
415 }
416
417 private recordSessionRelationship(notification: HarnessNotification): void {
418 if (notification.method !== 'subagent.started') return
419 const parentId = notification.params.parentSessionId
420 const childId = notification.params.childSessionId
421 if (typeof parentId === 'string' && parentId !== '' && typeof childId === 'string' && childId !== '' && parentId !== childId) {
422 this.sessionParents.set(childId, parentId)
423 }
424 }
425
426 private isDescendantOf(sessionId: string, rootSessionId: string): boolean {
427 const visited = new Set<string>()
428 let current = sessionId
429 while (!visited.has(current)) {
430 if (current === rootSessionId) return true
431 visited.add(current)
432 const parent = this.sessionParents.get(current)
433 if (parent === undefined) return false
434 current = parent
435 }
436 // The parent map only ever extends chains upward, so a cycle cannot form.
437 /* v8 ignore next */
438 return false
439 }
440
441 private failSubscriptions(error: Error): void {
442 for (const subscription of this.subscriptions.values()) subscription.fail(error)
443 }
444
445 private appendStderr(lines: string[]): void {
446 const kept = lines.filter(line => line.length > 0)
447 this.stderrTail.push(...kept)
448 if (this.stderrTail.length > STDERR_TAIL_LIMIT) {
449 this.stderrTail.splice(0, this.stderrTail.length - STDERR_TAIL_LIMIT)
450 }
451 }
452
453 private settleStreams(): Promise<void> {
454 return Promise.race([
455 this.streamsSettled,
456 new Promise<void>((resolve) => { setTimeout(resolve, STREAM_SETTLE_MS) }),
457 ])
458 }
459
460 private closedError(reason: string): TransportClosedError {
461 const parts = [`${this.runtime.description}: ${reason}`]
462 if (this.spawnError !== undefined) parts.push(`spawn error: ${this.spawnError.message}`)
463 if (this.exitCode !== undefined) parts.push(`exit code: ${String(this.exitCode)}`)
464 if (this.stderrTail.length > 0) parts.push(`stderr tail:\n${this.stderrTail.join('\n')}`)
465 return new TransportClosedError(parts.join('\n'))
466 }
467}
468
469/** Construct the transport against a generic process for package-local fake-runtime tests. */
470export function createProcessHarnessClient(options: RuntimeProcessOptions): HarnessClient {
471 const Constructor = HarnessClient as new (
472 publicOptions: HarnessClientOptions,
473 runtime: RuntimeProcessOptions,
474 ) => HarnessClient
475 return new Constructor({}, options)
476}
477
478/**
479 * Whether `value` is a plain JSON object (the wire-boundary shape probe).
480 * @param value - the wire value to probe.
481 * @returns `true` iff `value` is a non-null, non-array object.
482 */
483export function isRecord(value: unknown): value is Record<string, unknown> {
484 return typeof value === 'object' && value !== null && !Array.isArray(value)
485}
486
487/** The message of a thrown value (the transport only throws `Error`s; `String` covers the rest). */
488function errorMessage(error: unknown): string {
489 /* v8 ignore next -- the transport and dispose ladder reject only with Errors */
490 return error instanceof Error ? error.message : String(error)
491}