fix(sdk-client): address ds-review-bot findings

- api: resolve a relative workspace cwd to absolute before the handshake —
  the child spawns relative to the parent cwd, but the wire cwd is resolved
  again inside the child, so a relative value double-resolved
  (worker -> worker/worker).
- api: make the documented handshake retry real — HarnessClient.close() is
  permanent, so a failed initialize now reaps the runtime and swaps in a
  fresh client; DeepSeekHarness.close() is terminal and stops the respawns.
- api: validate session.event envelopes, assistant/message content, and
  session.finished reasons at the wire boundary — a malformed runtime
  surfaces as SdkProtocolError instead of type-invalid TurnResult data or a
  TypeError out of finalResponse.
- client: a throwing subscribe() filter fails and detaches only its own
  subscription (normalized to Error); sibling fan-out and the transport read
  loop are undisturbed.
- client: NotificationSubscription.close() drops its queued notifications,
  matching its documented contract; runtime-death fail() still leaves
  already-delivered items drainable.
- client: subscribe() after close()/runtime death returns a born-failed
  subscription so next() rejects instead of parking forever.
- client/transport: bounded requests abandon via AbortSignal — the transport
  drops the pending entry at timeout, so repeated bounded calls against a
  hung method retain no per-call state.

One test per finding; per-file coverage stays 100% on both packages.
This commit is contained in:
Tianyi Cui
2026-07-27 17:48:07 +08:00
parent 24d2384294
commit cf2b9e211d
9 changed files with 347 additions and 44 deletions

View File

@@ -8,9 +8,10 @@
*/
import { randomUUID } from 'node:crypto'
import { resolve } from 'node:path'
import type { SessionEvent, TurnEndReason } from '@deepseek-ai/dsh-session'
import { HarnessClient } from './client.ts'
import type { ContentBlock, DeepSeekHarnessOptions, HarnessNotification, TurnResult } from './types.ts'
import { HarnessClient, isRecord, SdkProtocolError } from './client.ts'
import type { ContentBlock, DeepSeekHarnessOptions, HarnessClientOptions, HarnessNotification, TurnResult } from './types.ts'
/**
* Reusable SDK for running DeepSeek Harness agent turns in a runtime
@@ -19,33 +20,52 @@ import type { ContentBlock, DeepSeekHarnessOptions, HarnessNotification, TurnRes
* child is reaped.
*/
export class DeepSeekHarness implements AsyncDisposable {
/** The underlying JSON-RPC client (exposed for low-level access). */
readonly client: HarnessClient
private clientInstance: HarnessClient
private readonly launch: HarnessClientOptions
private readonly cwd: string
private readonly provider: string
private readonly model: string
private initialized: Promise<void> | undefined
private closed = false
/** @param options - runtime launch spec plus the session route (cwd/provider/model). */
constructor(options: DeepSeekHarnessOptions) {
this.client = new HarnessClient(options.launch)
this.cwd = options.cwd ?? options.launch.cwd ?? process.cwd()
this.launch = options.launch
this.clientInstance = new HarnessClient(options.launch)
// Absolute before the handshake: the child spawns relative to THIS
// process's cwd, but the wire cwd is resolved again inside the child — a
// relative value would double-resolve (e.g. `worker` → `worker/worker`).
this.cwd = resolve(options.cwd ?? options.launch.cwd ?? process.cwd())
this.provider = options.provider ?? 'deepseek'
this.model = options.model ?? 'deepseek-v4-flash'
}
/**
* Start the subprocess and perform the `initialize` handshake once.
* The underlying JSON-RPC client (exposed for low-level access). A failed
* handshake reaps its runtime and swaps in a fresh instance, so do not
* cache this across a failed {@link start}.
* @returns the client currently owning the runtime subprocess.
*/
get client(): HarnessClient {
return this.clientInstance
}
/**
* Start the subprocess and perform the `initialize` handshake once. On
* failure the runtime is reaped and a fresh client replaces it
* (`HarnessClient.close` is permanent), so a later call retries with a new
* subprocess — unless {@link close} already ended this harness.
* @returns settlement of the (memoized) handshake.
*/
start(): Promise<void> {
this.initialized ??= (async () => {
try {
this.client.start()
await this.client.initialize({ cwd: this.cwd, provider: this.provider, model: this.model })
this.clientInstance.start()
await this.clientInstance.initialize({ cwd: this.cwd, provider: this.provider, model: this.model })
} catch (error) {
this.initialized = undefined
await this.client.close()
await this.clientInstance.close()
if (!this.closed) this.clientInstance = new HarnessClient(this.launch)
throw error
}
})()
@@ -73,11 +93,13 @@ export class DeepSeekHarness implements AsyncDisposable {
}
/**
* Shut down and reap the runtime subprocess. Idempotent.
* Shut down and reap the runtime subprocess. Idempotent and terminal —
* a closed harness no longer retries a failed handshake.
* @returns settlement of the complete teardown.
*/
close(): Promise<void> {
return this.client.close()
this.closed = true
return this.clientInstance.close()
}
/**
@@ -128,16 +150,26 @@ export class HarnessSession {
const subscription = client.subscribeSessionTree(this.id)
const collect = (notification: HarnessNotification): void => {
notifications.push(notification)
options?.onNotification?.(notification)
if (notification.method === 'session.event' && notification.params.sessionId === this.id) {
events.push(notification.params.event as SessionEvent)
// Wire boundary: the envelope feeds the typed TurnResult, so a
// malformed runtime surfaces as a protocol error, not as type-invalid
// data (or a TypeError out of finalResponse).
const event = validatedSessionEvent(notification.params.event)
notifications.push(notification)
options?.onNotification?.(notification)
events.push(event)
return
}
if (notification.method === 'session.finished' && notification.params.sessionId === this.id) {
reason = validatedTurnEndReason(notification.params.reason)
notifications.push(notification)
options?.onNotification?.(notification)
status = notification.params.status === 'ok' ? 'ok' : 'error'
reason = notification.params.reason as TurnEndReason | undefined
finished = true
return
}
notifications.push(notification)
options?.onNotification?.(notification)
}
const accepted = client.prompt(this.id, contentBlocks)
// Drain concurrently so observers see progress while the prompt request
@@ -175,6 +207,32 @@ export function normalizeInput(input: string | ContentBlock[]): ContentBlock[] {
return typeof input === 'string' ? [{ type: 'text', text: input }] : input
}
/** Validate a wire `session.event` envelope to the shape the typed result exposes. */
function validatedSessionEvent(value: unknown): SessionEvent {
if (!isRecord(value) || typeof value.type !== 'string') {
throw new SdkProtocolError(`session.event carried no event envelope: ${JSON.stringify(value)}`)
}
// The one variant this module reads into (finalResponse) must carry
// kind-tagged content blocks; other variants pass through under their
// envelope shape.
if (value.type === 'assistant/message') {
const content = isRecord(value.data) ? value.data.content : undefined
if (!Array.isArray(content) || !content.every(block => isRecord(block) && typeof block.type === 'string')) {
throw new SdkProtocolError(`assistant/message event carried malformed content: ${JSON.stringify(value)}`)
}
}
return value as unknown as SessionEvent
}
/** Validate a wire `session.finished` reason (absent, or a kind-tagged record). */
function validatedTurnEndReason(value: unknown): TurnEndReason | undefined {
if (value === undefined) return undefined
if (!isRecord(value) || typeof value.kind !== 'string') {
throw new SdkProtocolError(`session.finished carried a malformed reason: ${JSON.stringify(value)}`)
}
return value as unknown as TurnEndReason
}
/**
* Extract the concatenated text of the last assistant message.
* @param events - the turn's `session.event` payloads in wire order.

View File

@@ -81,8 +81,9 @@ export class NotificationSubscription implements AsyncIterable<HarnessNotificati
/**
* Await the next matching notification.
* @returns the notification; rejects once the runtime is closed or the
* subscription itself is closed.
* @returns the notification; after the runtime died, drains what was
* already delivered and then rejects; after {@link close}, rejects
* immediately (the queue is dropped).
*/
next(): Promise<HarnessNotification> {
const queued = this.state.queue.shift()
@@ -104,11 +105,15 @@ export class NotificationSubscription implements AsyncIterable<HarnessNotificati
/** Detach from the client; queued items drop and pending waiters reject. */
close(): void {
this.unsubscribe()
// The drop is part of this method's contract; a runtime-death fail() keeps
// the queue so already-delivered notifications remain drainable.
this.state.queue.length = 0
this.fail(new TransportClosedError('notification subscription closed'))
}
/**
* Reject pending and future waits (delivery stops; the first failure wins).
* Already-queued notifications remain drainable via {@link next}/{@link tryNext}.
* @param error - the terminal failure delivered to waiters.
*/
fail(error: Error): void {
@@ -117,11 +122,22 @@ export class NotificationSubscription implements AsyncIterable<HarnessNotificati
}
/**
* Deliver one notification to a waiter or the queue when the filter matches.
* Deliver one notification to a waiter or the queue when the filter
* matches. A throwing filter fails only THIS subscription (detached, the
* throw becomes its terminal error) — it never disturbs sibling
* subscriptions or the transport's read loop, mirroring the Python client.
* @param notification - the wire notification to deliver.
*/
push(notification: HarnessNotification): void {
if (this.state.filter !== undefined && !this.state.filter(notification)) return
let matches: boolean
try {
matches = this.state.filter === undefined || this.state.filter(notification)
} catch (error) {
this.unsubscribe()
this.fail(error instanceof Error ? error : new Error(String(error)))
return
}
if (!matches) return
const waiter = this.state.waiters.shift()
if (waiter !== undefined) waiter.resolve(notification)
else this.state.queue.push(notification)
@@ -273,22 +289,18 @@ export class HarnessClient {
const transport = this.transport
/* v8 ignore next -- start() either sets the transport or throws */
if (transport === undefined) throw new TransportClosedError('DeepSeek Harness runtime is not running')
const pending = transport.request(method, params ?? {})
const timeout = timeoutMs ?? this.options.requestTimeoutMs
try {
if (timeout === undefined) return await pending
let timer: NodeJS.Timeout | undefined
if (timeout === undefined) return await transport.request(method, params ?? {})
// The abort signal makes the timeout an abandonment: the transport drops
// its pending entry, so repeated bounded requests against a hung method
// retain no per-call state (the server-side work still runs to close).
const abandon = new AbortController()
const timer = setTimeout(() => {
abandon.abort(new RequestTimeoutError(`${method} timed out after ${timeout}ms waiting for the DeepSeek Harness runtime`))
}, timeout)
try {
return await Promise.race([
pending,
new Promise<never>((_, reject) => {
timer = setTimeout(() => {
// The abandoned wire promise settles on close; keep it handled.
pending.catch(() => {})
reject(new RequestTimeoutError(`${method} timed out after ${timeout}ms waiting for the DeepSeek Harness runtime`))
}, timeout)
}),
])
return await transport.request(method, params ?? {}, abandon.signal)
} finally {
clearTimeout(timer)
}
@@ -303,12 +315,18 @@ export class HarnessClient {
/**
* Subscribe to server notifications.
* @param filter - optional predicate; omitted means every notification.
* @returns the subscription handle; close it to stop delivery.
* @returns the subscription handle; close it to stop delivery. After
* {@link close} or runtime death the handle is born failed — there is no
* producer left, so `next()` rejects instead of waiting forever.
*/
subscribe(filter?: NotificationFilter): NotificationSubscription {
const id = String(this.subscriptionSerial++)
const state: SubscriptionState = { queue: [], waiters: [], filter, failure: undefined }
const subscription = new NotificationSubscription(state, () => { this.subscriptions.delete(id) })
if (this.closeTask !== undefined || this.exitCode !== undefined || this.spawnError !== undefined) {
subscription.fail(this.closedError('DeepSeek Harness runtime closed'))
return subscription
}
this.subscriptions.set(id, subscription)
return subscription
}