The native adapter's route was named deepseek, colliding with pi-ai's catalog provider of the same name, so the two DeepSeek paths could never be mounted side by side. The web settings page needs both configurable at once. Compositions, fixtures, goldens, scaffolding defaults, and docs all move together (pre-release, no shim); TUI/session-query-spill/ missing-credential goldens re-recorded through their keyless refresh modes because provider-name length shifts box padding and spill truncation points.
260 lines
10 KiB
TypeScript
260 lines
10 KiB
TypeScript
/**
|
|
* High-level turns API over {@link HarnessClient}: `DeepSeekHarness` owns one
|
|
* runtime subprocess across many sessions; `HarnessSession.run` sends a
|
|
* prompt and settles with the final response once `session.finished` arrives.
|
|
* Mirrors the Python SDK's `DeepSeekHarness`/`Session` pair.
|
|
*
|
|
* @module @deepseek-ai/dsh-sdk-client/api
|
|
*/
|
|
|
|
import { randomUUID } from 'node:crypto'
|
|
import { resolve } from 'node:path'
|
|
import type { SessionEvent, TurnEndReason } from '@deepseek-ai/dsh-session'
|
|
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
|
|
* subprocess. The subprocess starts lazily on first use and stays owned by
|
|
* this instance until {@link close}; always close (or `await using`) so the
|
|
* child is reaped.
|
|
*/
|
|
export class DeepSeekHarness implements AsyncDisposable {
|
|
private clientInstance: HarnessClient
|
|
private readonly launch: HarnessClientOptions
|
|
private readonly cwd: string
|
|
private readonly provider: string
|
|
private readonly model: string
|
|
private readonly maxTokens: number | undefined
|
|
private initialized: Promise<void> | undefined
|
|
private closed = false
|
|
|
|
/** @param options - runtime launch spec plus the session route (cwd/provider/model). */
|
|
constructor(options: DeepSeekHarnessOptions) {
|
|
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-official'
|
|
this.model = options.model ?? 'deepseek-v4-flash'
|
|
this.maxTokens = options.maxTokens
|
|
}
|
|
|
|
/**
|
|
* 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.clientInstance.start()
|
|
await this.clientInstance.initialize({
|
|
cwd: this.cwd,
|
|
provider: this.provider,
|
|
model: this.model,
|
|
...this.maxTokens === undefined ? {} : { maxTokens: this.maxTokens },
|
|
})
|
|
} catch (error) {
|
|
this.initialized = undefined
|
|
await this.clientInstance.close()
|
|
if (!this.closed) this.clientInstance = new HarnessClient(this.launch)
|
|
throw error
|
|
}
|
|
})()
|
|
return this.initialized
|
|
}
|
|
|
|
/**
|
|
* Open a session handle (no wire traffic; the runtime creates the session
|
|
* on its first prompt).
|
|
* @param sessionId - explicit id to reuse; omitted mints a fresh one.
|
|
* @returns the session handle.
|
|
*/
|
|
session(sessionId?: string): HarnessSession {
|
|
return new HarnessSession(this, sessionId ?? `session-${randomUUID().replaceAll('-', '')}`)
|
|
}
|
|
|
|
/**
|
|
* Run one prompt on a fresh (or named) session.
|
|
* @param input - prompt text, or content blocks sent verbatim.
|
|
* @param options - optional session id and per-notification observer.
|
|
* @returns the settled turn result.
|
|
*/
|
|
run(input: string | ContentBlock[], options?: RunOptions): Promise<TurnResult> {
|
|
return this.session(options?.sessionId).run(input, options)
|
|
}
|
|
|
|
/**
|
|
* 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> {
|
|
this.closed = true
|
|
return this.clientInstance.close()
|
|
}
|
|
|
|
/**
|
|
* `await using` support: {@link close}.
|
|
* @returns settlement of the teardown.
|
|
*/
|
|
[Symbol.asyncDispose](): Promise<void> {
|
|
return this.close()
|
|
}
|
|
}
|
|
|
|
/** Per-run options: target session and streaming observer. */
|
|
export interface RunOptions {
|
|
/** Session id to run on; omitted mints a fresh session per call. */
|
|
sessionId?: string
|
|
/** Observer invoked with every notification for this session tree, in wire order. */
|
|
onNotification?: (notification: HarnessNotification) => void
|
|
}
|
|
|
|
/**
|
|
* One SDK session: a stable id plus the turn loop that pairs a
|
|
* `session/prompt` with its `session.finished`.
|
|
*/
|
|
export class HarnessSession {
|
|
/**
|
|
* @param harness - the owning harness (supplies the client and handshake).
|
|
* @param id - the wire session id this handle runs on.
|
|
*/
|
|
constructor(readonly harness: DeepSeekHarness, readonly id: string) {}
|
|
|
|
/**
|
|
* Run one prompt turn to settlement.
|
|
* @param input - prompt text, or content blocks sent verbatim.
|
|
* @param options - optional per-notification observer.
|
|
* @returns the settled turn result; rejects on transport loss, timeout, or
|
|
* a protocol error — never on a model-level failure (that is
|
|
* `status: 'error'` in the result).
|
|
*/
|
|
async run(input: string | ContentBlock[], options?: Pick<RunOptions, 'onNotification'>): Promise<TurnResult> {
|
|
await this.harness.start()
|
|
const client = this.harness.client
|
|
const contentBlocks = normalizeInput(input)
|
|
const events: SessionEvent[] = []
|
|
const notifications: HarnessNotification[] = []
|
|
let status: TurnResult['status'] = 'error'
|
|
let reason: TurnEndReason | undefined
|
|
let finished = false
|
|
|
|
const subscription = client.subscribeSessionTree(this.id)
|
|
const collect = (notification: HarnessNotification): void => {
|
|
if (notification.method === 'session.event' && notification.params.sessionId === this.id) {
|
|
// 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'
|
|
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
|
|
// is still pending (its response arrives only after settlement).
|
|
const drain = (async () => {
|
|
while (!finished) collect(await subscription.next())
|
|
})()
|
|
try {
|
|
await Promise.all([accepted, drain])
|
|
} finally {
|
|
// On a prompt rejection the drain is still parked on next(); closing the
|
|
// subscription settles it, and the swallow keeps that secondary
|
|
// TransportClosedError from surfacing as an unhandled rejection.
|
|
subscription.close()
|
|
await drain.catch(() => {})
|
|
}
|
|
|
|
return {
|
|
sessionId: this.id,
|
|
status,
|
|
reason,
|
|
finalResponse: finalResponse(events),
|
|
events,
|
|
notifications,
|
|
}
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Normalize run input: a string becomes one text block; blocks pass verbatim.
|
|
* @param input - prompt text or content blocks.
|
|
* @returns the content blocks to send.
|
|
*/
|
|
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 message = isRecord(value.data) ? value.data.message : undefined
|
|
const content = isRecord(message) ? message.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.
|
|
* @returns the final response text, or `''` when no assistant message exists.
|
|
*/
|
|
export function finalResponse(events: SessionEvent[]): string {
|
|
for (let index = events.length - 1; index >= 0; index--) {
|
|
const event = events[index]
|
|
if (event?.type !== 'assistant/message') continue
|
|
return event.data.message.content
|
|
.filter((block): block is ContentBlock & { type: 'text' } => block.type === 'text')
|
|
.map(block => block.text)
|
|
.join('')
|
|
}
|
|
return ''
|
|
}
|