1
/**2
* High-level run API over {@link HarnessClient}: `DeepSeekHarness` owns one3
* runtime subprocess across many sessions; `HarnessSession.run` sends a4
* prompt and settles when the whole agent next becomes idle.5
*6
* @module @deepseek-ai/dsh-sdk-client/api7
*/9
import { randomUUID } from 'node:crypto'10
import { resolve } from 'node:path'11
import type { SessionEvent, TurnEndReason } from '@deepseek-ai/dsh-session'12
import { createProcessHarnessClient, HarnessClient, isRecord, SdkProtocolError } from './client.ts'13
import type { RuntimeProcessOptions } from './launch.ts'14
import type { ContentBlock, DeepSeekHarnessOptions, HarnessNotification, RunResult, SdkPromptContentBlock } from './types.ts'16
/**17
* Reusable SDK for running DeepSeek Harness agent turns in a runtime18
* subprocess. The subprocess starts lazily on first use and stays owned by19
* this instance until {@link close}; always close (or `await using`) so the20
* child is reaped.21
*/22
export class DeepSeekHarness implements AsyncDisposable {23
private clientInstance: HarnessClient24
private readonly createClient: () => HarnessClient25
private readonly cwd: string26
private readonly provider: string27
private readonly model: string28
private readonly reasoningEffort: DeepSeekHarnessOptions['reasoningEffort']29
private readonly maxTokens: number | undefined30
private initialized: Promise<void> | undefined31
private closed = false33
/** @param options - dsh launch configuration plus the session route, effort, and output cap. */34
constructor(options?: DeepSeekHarnessOptions)35
constructor(options: DeepSeekHarnessOptions = {}, clientFactory?: () => HarnessClient) {36
this.createClient = clientFactory ?? (() => new HarnessClient(options))37
this.clientInstance = this.createClient()38
// Absolute before the handshake: the child spawns relative to THIS39
// process's cwd, but the wire cwd is resolved again inside the child — a40
// relative value would double-resolve (e.g. `worker` → `worker/worker`).41
this.cwd = resolve(options.cwd ?? options.processCwd ?? process.cwd())42
this.provider = options.provider ?? 'deepseek-official'43
this.model = options.model ?? 'deepseek-v4-flash'44
this.reasoningEffort = options.reasoningEffort45
this.maxTokens = options.maxTokens46
}48
/**49
* The underlying JSON-RPC client (exposed for low-level access). A failed50
* handshake swaps in a fresh instance only after cleanup proves the runtime51
* exited; cleanup failure retains this client, so do not cache it across a52
* failed {@link start}.53
* @returns the client currently owning the runtime subprocess.54
*/55
get client(): HarnessClient {56
return this.clientInstance57
}59
/**60
* Start the subprocess and perform the `initialize` handshake once. On61
* failure, successful SDK-owned cleanup reaps the runtime and installs a62
* fresh client (`HarnessClient.close` is permanent), so a later call retries63
* with a new subprocess unless {@link close} already ended this harness. If64
* cleanup also fails, rejects with an `AggregateError` whose ordered errors65
* preserve both causes and retains the failed client rather than spawning66
* alongside a process whose exit was not proved.67
* @returns settlement of the (memoized) handshake.68
*/69
start(): Promise<void> {70
this.initialized ??= (async () => {71
try {72
this.clientInstance.start()73
await this.clientInstance.initialize({74
cwd: this.cwd,75
provider: this.provider,76
model: this.model,77
...this.reasoningEffort === undefined ? {} : { reasoningEffort: this.reasoningEffort },78
...this.maxTokens === undefined ? {} : { maxTokens: this.maxTokens },79
})80
} catch (error) {81
this.initialized = undefined82
try {83
await this.clientInstance.close()84
} catch (cleanupError: unknown) {85
throw new AggregateError(86
[error, cleanupError],87
'DeepSeek Harness initialization and cleanup failed',88
)89
}90
if (!this.closed) this.clientInstance = this.createClient()91
throw error92
}93
})()94
return this.initialized95
}97
/**98
* Open a session handle (no wire traffic; the runtime creates the session99
* on its first prompt).100
* @param sessionId - explicit id to reuse; omitted mints a fresh one.101
* @returns the session handle.102
*/103
session(sessionId?: string): HarnessSession {104
return new HarnessSession(this, sessionId ?? `session-${randomUUID().replaceAll('-', '')}`)105
}107
/**108
* Run one prompt on a fresh (or named) session.109
* @param input - prompt text, or content blocks sent verbatim.110
* @param options - optional session id and per-notification observer.111
* @returns the owned activity interval.112
*/113
run(input: string | SdkPromptContentBlock[], options?: RunOptions): Promise<RunResult> {114
return this.session(options?.sessionId).run(input, options)115
}117
/**118
* Shut down and reap the runtime subprocess. Idempotent and terminal —119
* a closed harness no longer retries a failed handshake.120
* @returns settlement of the complete teardown.121
*/122
close(): Promise<void> {123
this.closed = true124
return this.clientInstance.close()125
}127
/**128
* `await using` support: {@link close}.129
* @returns settlement of the teardown.130
*/131
[Symbol.asyncDispose](): Promise<void> {132
return this.close()133
}134
}136
/** Construct the high-level API against a generic process for package-local fake-runtime tests. */137
export function createProcessDeepSeekHarness(138
runtime: RuntimeProcessOptions,139
options: DeepSeekHarnessOptions = {},140
): DeepSeekHarness {141
const Constructor = DeepSeekHarness as new (142
publicOptions: DeepSeekHarnessOptions,143
clientFactory: () => HarnessClient,144
) => DeepSeekHarness145
return new Constructor({146
...runtime.cwd === undefined ? {} : { processCwd: runtime.cwd },147
...options,148
}, () => createProcessHarnessClient(runtime))149
}151
/** Per-run options: target session and streaming observer. */152
export interface RunOptions {153
/** Session id to run on; omitted mints a fresh session per call. */154
sessionId?: string155
/** Observer invoked with every notification for this session tree, in wire order. */156
onNotification?: (notification: HarnessNotification) => void157
}159
/**160
* One SDK session: a stable id plus owned activity intervals.161
*/162
export class HarnessSession {163
/**164
* @param harness - the owning harness (supplies the client and handshake).165
* @param id - the wire session id this handle runs on.166
*/167
constructor(readonly harness: DeepSeekHarness, readonly id: string) {}169
/**170
* Queue one prompt, then observe the whole session through its next idle.171
* @param input - prompt text, or content blocks sent verbatim.172
* @param options - optional per-notification observer.173
* @returns the owned activity interval; rejects on transport loss, timeout,174
* or a protocol error.175
*/176
async run(input: string | SdkPromptContentBlock[], options?: Pick<RunOptions, 'onNotification'>): Promise<RunResult> {177
await this.harness.start()178
const client = this.harness.client179
const contentBlocks = normalizeInput(input)180
const events: SessionEvent[] = []181
const notifications: HarnessNotification[] = []183
const subscription = client.subscribeSessionTree(this.id)184
const collect = (notification: HarnessNotification): void => {185
if (notification.method === 'session.event' && notification.params.sessionId === this.id) {186
// Wire boundary: the envelope feeds the typed RunResult, so a187
// malformed runtime surfaces as a protocol error, not as type-invalid188
// data (or a TypeError out of finalResponse).189
const event = validatedSessionEvent(notification.params.event)190
notifications.push(notification)191
options?.onNotification?.(notification)192
events.push(event)193
return194
}195
notifications.push(notification)196
options?.onNotification?.(notification)197
}198
try {199
const messageId = await client.prompt(this.id, contentBlocks)200
let received = false201
while (true) {202
const notification = await subscription.next()203
if (!received) {204
if (notification.method !== 'session.event'205
|| notification.params.sessionId !== this.id206
|| !isInboxReceipt(notification.params.event, messageId)) continue207
received = true208
}209
collect(notification)210
if (notification.method === 'session.status'211
&& notification.params.sessionId === this.id212
&& notification.params.status === 'idle') break213
}214
} finally {215
subscription.close()216
}218
return {219
sessionId: this.id,220
finalResponse: finalResponse(events),221
events,222
notifications,223
}224
}225
}227
/**228
* Normalize run input: a string becomes one text block; blocks pass verbatim.229
* @param input - prompt text or content blocks.230
* @returns the content blocks to send.231
*/232
export function normalizeInput(input: string | SdkPromptContentBlock[]): SdkPromptContentBlock[] {233
return typeof input === 'string' ? [{ type: 'text', text: input }] : input234
}236
/** Validate the provider-read fields of one wire turn-end reason. */237
function validatedTurnEndReason(value: unknown): TurnEndReason {238
if (!isRecord(value) || typeof value.kind !== 'string') {239
throw new SdkProtocolError(`turn/end carried no reason envelope: ${JSON.stringify(value)}`)240
}241
if (value.kind === 'aborted') {242
if (!isRecord(value.reason) || typeof value.reason.kind !== 'string') {243
throw new SdkProtocolError(`turn/end carried a malformed aborted reason: ${JSON.stringify(value)}`)244
}245
switch (value.reason.kind) {246
case 'user':247
case 'parent':248
case 'disposed':249
case 'legacy':250
break251
case 'hook':252
if (typeof value.reason.reason !== 'string') {253
throw new SdkProtocolError(`turn/end carried a malformed hook abort reason: ${JSON.stringify(value)}`)254
}255
break256
default:257
throw new SdkProtocolError(`turn/end carried an unknown abort reason: ${JSON.stringify(value)}`)258
}259
}260
return value as TurnEndReason261
}263
/** Validate the fields in a wire `session.event` envelope before returning the typed result. */264
function validatedSessionEvent(value: unknown): SessionEvent {265
if (!isRecord(value) || typeof value.type !== 'string') {266
throw new SdkProtocolError(`session.event carried no event envelope: ${JSON.stringify(value)}`)267
}268
// The one variant this module reads into (finalResponse) must carry269
// kind-tagged content blocks; other variants pass through under their270
// envelope shape.271
if (value.type === 'assistant/message') {272
const message = isRecord(value.data) ? value.data.message : undefined273
const content = isRecord(message) ? message.content : undefined274
if (!Array.isArray(content) || !content.every(block => isRecord(block) && typeof block.type === 'string')) {275
throw new SdkProtocolError(`assistant/message event carried malformed content: ${JSON.stringify(value)}`)276
}277
}278
if (value.type === 'turn/end') {279
const data = isRecord(value.data) ? value.data : undefined280
if (data === undefined) {281
throw new SdkProtocolError(`turn/end event carried malformed data: ${JSON.stringify(value)}`)282
}283
validatedTurnEndReason(data.reason)284
}285
return value as SessionEvent286
}288
/** Whether a raw session event is the durable enqueue receipt for `messageId`. */289
function isInboxReceipt(value: unknown, messageId: string): boolean {290
if (!isRecord(value) || value.type !== 'agent/inbox/spliced' || !isRecord(value.data)) return false291
const inserted = value.data.inserted292
return Array.isArray(inserted) && inserted.some(message => isRecord(message) && message.id === messageId)293
}295
/**296
* Extract the concatenated text of the last assistant message.297
* @param events - the activity interval's `session.event` payloads in wire order.298
* @returns the final response text, or `''` when no assistant message exists.299
*/300
export function finalResponse(events: SessionEvent[]): string {301
for (let index = events.length - 1; index >= 0; index--) {302
const event = events[index]303
if (event?.type !== 'assistant/message') continue304
return event.data.message.content305
.filter((block): block is ContentBlock & { type: 'text' } => block.type === 'text')306
.map(block => block.text)307
.join('')308
}309
return ''310
}