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, speaks4
* the `@deepseek-ai/dsh-sdk-protocol` wire over the child's stdio, fans5
* server notifications out to subscriptions, and tears the child down to6
* quiescence through a private EOF → SIGTERM → SIGKILL ladder. The design7
* twin is the Python SDK's `HarnessClient` (`python/sdk`); both drive the8
* same runtime protocol. This client runs OUTSIDE any harness context, so it9
* spawns directly rather than through the `dsh-subprocess` service — the10
* seam's documented exception for SDK-managed transports.11
*12
* @module @deepseek-ai/dsh-sdk-client/client13
*/15
import { spawn, type ChildProcess } from 'node:child_process'16
import {17
JsonRpcLineTransport,18
JsonRpcResponseError,19
type InitializeParams,20
type InitializeResult,21
type SessionPromptParams,22
type SdkPromptContentBlock,23
} from '@deepseek-ai/dsh-sdk-protocol'24
import { disposeRuntimeProcess } from './dispose.ts'25
import { resolveDshLaunch, type RuntimeProcessOptions } from './launch.ts'26
import type { HarnessClientOptions, HarnessNotification, NotificationFilter } from './types.ts'28
/** Retained stderr lines used to diagnose an unexpected runtime death. */29
const STDERR_TAIL_LIMIT = 40031
/** Grace for the runtime's stdio streams to settle after its exit edge. */32
const STREAM_SETTLE_MS = 10034
/**35
* The runtime subprocess is gone or unusable: it exited, its stdio closed, or36
* it was never launchable. The message carries the exit code and a stderr37
* tail when available.38
*/39
export 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
}47
/** A request exceeded {@link HarnessClientOptions.requestTimeoutMs}. */48
export class RequestTimeoutError extends Error {49
/** @param message - which method timed out. */50
constructor(message: string) {51
super(message)52
this.name = 'RequestTimeoutError'53
}54
}56
/**57
* The runtime answered outside its documented protocol (for example a58
* `session/prompt` response without `accepted: true`).59
*/60
export class SdkProtocolError extends Error {61
/** @param message - the protocol violation description. */62
constructor(message: string) {63
super(message)64
this.name = 'SdkProtocolError'65
}66
}68
interface SubscriptionState {69
readonly queue: HarnessNotification[]70
readonly waiters: { resolve: (item: HarnessNotification) => void; reject: (error: Error) => void }[]71
readonly filter: NotificationFilter | undefined72
failure: Error | undefined73
}75
/** One client-side notification stream returned by {@link HarnessClient.subscribe}. */76
export interface NotificationSubscription extends AsyncIterable<HarnessNotification> {77
/**78
* Await the next matching notification.79
* @returns the notification; after the runtime died, drains what was80
* already delivered and then rejects; after {@link close}, rejects81
* immediately (the queue is dropped).82
*/83
next(): Promise<HarnessNotification>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 | undefined91
/** Detach from the client; queued items drop and pending waiters reject. */92
close(): void93
}95
/** Internal producer side of a public notification subscription. */96
class NotificationSubscriptionImpl implements NotificationSubscription {97
constructor(98
private readonly state: SubscriptionState,99
private readonly unsubscribe: () => void,100
) {}102
/**103
* Await the next matching notification.104
* @returns the notification; after the runtime died, drains what was105
* already delivered and then rejects; after {@link close}, rejects106
* 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
}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
}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() keeps129
// the queue so already-delivered notifications remain drainable.130
this.state.queue.length = 0131
this.fail(new TransportClosedError('notification subscription closed'))132
}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 ??= error141
for (const waiter of this.state.waiters.splice(0)) waiter.reject(this.state.failure)142
}144
/**145
* Deliver one notification to a waiter or the queue when the filter146
* matches. A throwing filter fails only THIS subscription (detached, the147
* throw becomes its terminal error) — it never disturbs sibling148
* subscriptions or the transport's read loop.149
* @param notification - the wire notification to deliver.150
*/151
push(notification: HarnessNotification): void {152
let matches: boolean153
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
return159
}160
if (!matches) return161
const waiter = this.state.waiters.shift()162
if (waiter !== undefined) waiter.resolve(notification)163
else this.state.queue.push(notification)164
}166
/**167
* Iterate notifications until the subscription or runtime closes (the168
* 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
}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 instance180
* until {@link close}, which requests protocol `shutdown` and then walks the181
* shared EOF → SIGTERM → SIGKILL dispose ladder to quiescence. There is no182
* wire-level cancel: a timed-out request stays running server-side until the183
* runtime is closed.184
*/185
export class HarnessClient {186
/** Original public dsh launch and timeout options for this client. */187
readonly options: HarnessClientOptions188
private readonly runtime: RuntimeProcessOptions189
private child: ChildProcess | undefined190
private transport: JsonRpcLineTransport | undefined191
private readonly stderrTail: string[] = []192
private readonly subscriptions = new Map<string, NotificationSubscriptionImpl>()193
private readonly sessionParents = new Map<string, string>()194
private subscriptionSerial = 0195
private exitCode: number | null | undefined196
private spawnError: Error | undefined197
private streamsSettled: Promise<void> = Promise.resolve()198
private closeTask: Promise<void> | undefined200
/** @param options - dsh profile, patch, home, process, environment, and timeout options. */201
constructor(options?: HarnessClientOptions)202
constructor(options: HarnessClientOptions = {}, runtime?: RuntimeProcessOptions) {203
this.options = options204
this.runtime = runtime ?? resolveDshLaunch(options)205
}207
/**208
* Spawn the runtime subprocess and start reading frames. Idempotent while209
* 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) return214
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 = child220
child.once('error', (error) => {221
this.spawnError = error222
// A spawn failure destroys the pipes without an input 'end' edge, so the223
// 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 is228
// 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 += chunk236
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!: () => void243
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 = true251
maybeSettle()252
})253
child.once('exit', (code) => {254
this.exitCode = code255
settled.exited = true256
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 that262
// 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 = transport269
}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
}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.messageId298
}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 a306
* protocol error response, {@link RequestTimeoutError} on timeout, and307
* {@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 of312
// 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.transport318
/* 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.requestTimeoutMs321
try {322
if (timeout === undefined) return await transport.request(method, params ?? {})323
// The abort signal makes the timeout an abandonment: the transport drops324
// its pending entry, so repeated bounded requests against a hung method325
// 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 error338
// Transport-level failures gain process context: exit code + stderr tail.339
await this.settleStreams()340
throw this.closedError(errorMessage(error))341
}342
}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. After348
* {@link close} or runtime death the handle is born failed — there is no349
* 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 subscription358
}359
this.subscriptions.set(id, subscription)360
return subscription361
}363
/**364
* Subscribe to one session and the descendants discovered from365
* `subagent.started` lineage edges. The runtime notifies for every session366
* 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.params373
if (notification.method === 'subagent.started' || notification.method === 'subagent.finished') {374
const parentId = params.parentSessionId375
if (typeof parentId === 'string' && this.isDescendantOf(parentId, sessionId)) return true376
return params.childSessionId === sessionId377
}378
const relatedId = params.sessionId379
return typeof relatedId === 'string' && this.isDescendantOf(relatedId, sessionId)380
})381
}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.closeTask392
}394
private async performClose(): Promise<void> {395
const child = this.child396
if (child === undefined) return397
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 teardown401
// 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
}412
private dispatchNotification(notification: HarnessNotification): void {413
this.recordSessionRelationship(notification)414
for (const subscription of this.subscriptions.values()) subscription.push(notification)415
}417
private recordSessionRelationship(notification: HarnessNotification): void {418
if (notification.method !== 'subagent.started') return419
const parentId = notification.params.parentSessionId420
const childId = notification.params.childSessionId421
if (typeof parentId === 'string' && parentId !== '' && typeof childId === 'string' && childId !== '' && parentId !== childId) {422
this.sessionParents.set(childId, parentId)423
}424
}426
private isDescendantOf(sessionId: string, rootSessionId: string): boolean {427
const visited = new Set<string>()428
let current = sessionId429
while (!visited.has(current)) {430
if (current === rootSessionId) return true431
visited.add(current)432
const parent = this.sessionParents.get(current)433
if (parent === undefined) return false434
current = parent435
}436
// The parent map only ever extends chains upward, so a cycle cannot form.437
/* v8 ignore next */438
return false439
}441
private failSubscriptions(error: Error): void {442
for (const subscription of this.subscriptions.values()) subscription.fail(error)443
}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
}453
private settleStreams(): Promise<void> {454
return Promise.race([455
this.streamsSettled,456
new Promise<void>((resolve) => { setTimeout(resolve, STREAM_SETTLE_MS) }),457
])458
}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
}469
/** Construct the transport against a generic process for package-local fake-runtime tests. */470
export function createProcessHarnessClient(options: RuntimeProcessOptions): HarnessClient {471
const Constructor = HarnessClient as new (472
publicOptions: HarnessClientOptions,473
runtime: RuntimeProcessOptions,474
) => HarnessClient475
return new Constructor({}, options)476
}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
*/483
export function isRecord(value: unknown): value is Record<string, unknown> {484
return typeof value === 'object' && value !== null && !Array.isArray(value)485
}487
/** The message of a thrown value (the transport only throws `Error`s; `String` covers the rest). */488
function 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
}