返回源码地图

packages/sdk/client/src/api.ts

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

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

1/**
2 * High-level run API over {@link HarnessClient}: `DeepSeekHarness` owns one
3 * runtime subprocess across many sessions; `HarnessSession.run` sends a
4 * prompt and settles when the whole agent next becomes idle.
5 *
6 * @module @deepseek-ai/dsh-sdk-client/api
7 */
8
9import { randomUUID } from 'node:crypto'
10import { resolve } from 'node:path'
11import type { SessionEvent, TurnEndReason } from '@deepseek-ai/dsh-session'
12import { createProcessHarnessClient, HarnessClient, isRecord, SdkProtocolError } from './client.ts'
13import type { RuntimeProcessOptions } from './launch.ts'
14import type { ContentBlock, DeepSeekHarnessOptions, HarnessNotification, RunResult, SdkPromptContentBlock } from './types.ts'
15
16/**
17 * Reusable SDK for running DeepSeek Harness agent turns in a runtime
18 * subprocess. The subprocess starts lazily on first use and stays owned by
19 * this instance until {@link close}; always close (or `await using`) so the
20 * child is reaped.
21 */
22export class DeepSeekHarness implements AsyncDisposable {
23 private clientInstance: HarnessClient
24 private readonly createClient: () => HarnessClient
25 private readonly cwd: string
26 private readonly provider: string
27 private readonly model: string
28 private readonly reasoningEffort: DeepSeekHarnessOptions['reasoningEffort']
29 private readonly maxTokens: number | undefined
30 private initialized: Promise<void> | undefined
31 private closed = false
32
33 /** @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 THIS
39 // process's cwd, but the wire cwd is resolved again inside the child — a
40 // 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.reasoningEffort
45 this.maxTokens = options.maxTokens
46 }
47
48 /**
49 * The underlying JSON-RPC client (exposed for low-level access). A failed
50 * handshake swaps in a fresh instance only after cleanup proves the runtime
51 * exited; cleanup failure retains this client, so do not cache it across a
52 * failed {@link start}.
53 * @returns the client currently owning the runtime subprocess.
54 */
55 get client(): HarnessClient {
56 return this.clientInstance
57 }
58
59 /**
60 * Start the subprocess and perform the `initialize` handshake once. On
61 * failure, successful SDK-owned cleanup reaps the runtime and installs a
62 * fresh client (`HarnessClient.close` is permanent), so a later call retries
63 * with a new subprocess unless {@link close} already ended this harness. If
64 * cleanup also fails, rejects with an `AggregateError` whose ordered errors
65 * preserve both causes and retains the failed client rather than spawning
66 * 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 = undefined
82 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 error
92 }
93 })()
94 return this.initialized
95 }
96
97 /**
98 * Open a session handle (no wire traffic; the runtime creates the session
99 * 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 }
106
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 }
116
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 = true
124 return this.clientInstance.close()
125 }
126
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}
135
136/** Construct the high-level API against a generic process for package-local fake-runtime tests. */
137export function createProcessDeepSeekHarness(
138 runtime: RuntimeProcessOptions,
139 options: DeepSeekHarnessOptions = {},
140): DeepSeekHarness {
141 const Constructor = DeepSeekHarness as new (
142 publicOptions: DeepSeekHarnessOptions,
143 clientFactory: () => HarnessClient,
144 ) => DeepSeekHarness
145 return new Constructor({
146 ...runtime.cwd === undefined ? {} : { processCwd: runtime.cwd },
147 ...options,
148 }, () => createProcessHarnessClient(runtime))
149}
150
151/** Per-run options: target session and streaming observer. */
152export interface RunOptions {
153 /** Session id to run on; omitted mints a fresh session per call. */
154 sessionId?: string
155 /** Observer invoked with every notification for this session tree, in wire order. */
156 onNotification?: (notification: HarnessNotification) => void
157}
158
159/**
160 * One SDK session: a stable id plus owned activity intervals.
161 */
162export 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) {}
168
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.client
179 const contentBlocks = normalizeInput(input)
180 const events: SessionEvent[] = []
181 const notifications: HarnessNotification[] = []
182
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 a
187 // malformed runtime surfaces as a protocol error, not as type-invalid
188 // 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 return
194 }
195 notifications.push(notification)
196 options?.onNotification?.(notification)
197 }
198 try {
199 const messageId = await client.prompt(this.id, contentBlocks)
200 let received = false
201 while (true) {
202 const notification = await subscription.next()
203 if (!received) {
204 if (notification.method !== 'session.event'
205 || notification.params.sessionId !== this.id
206 || !isInboxReceipt(notification.params.event, messageId)) continue
207 received = true
208 }
209 collect(notification)
210 if (notification.method === 'session.status'
211 && notification.params.sessionId === this.id
212 && notification.params.status === 'idle') break
213 }
214 } finally {
215 subscription.close()
216 }
217
218 return {
219 sessionId: this.id,
220 finalResponse: finalResponse(events),
221 events,
222 notifications,
223 }
224 }
225}
226
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 */
232export function normalizeInput(input: string | SdkPromptContentBlock[]): SdkPromptContentBlock[] {
233 return typeof input === 'string' ? [{ type: 'text', text: input }] : input
234}
235
236/** Validate the provider-read fields of one wire turn-end reason. */
237function 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 break
251 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 break
256 default:
257 throw new SdkProtocolError(`turn/end carried an unknown abort reason: ${JSON.stringify(value)}`)
258 }
259 }
260 return value as TurnEndReason
261}
262
263/** Validate the fields in a wire `session.event` envelope before returning the typed result. */
264function 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 carry
269 // kind-tagged content blocks; other variants pass through under their
270 // envelope shape.
271 if (value.type === 'assistant/message') {
272 const message = isRecord(value.data) ? value.data.message : undefined
273 const content = isRecord(message) ? message.content : undefined
274 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 : undefined
280 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 SessionEvent
286}
287
288/** Whether a raw session event is the durable enqueue receipt for `messageId`. */
289function isInboxReceipt(value: unknown, messageId: string): boolean {
290 if (!isRecord(value) || value.type !== 'agent/inbox/spliced' || !isRecord(value.data)) return false
291 const inserted = value.data.inserted
292 return Array.isArray(inserted) && inserted.some(message => isRecord(message) && message.id === messageId)
293}
294
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 */
300export 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') continue
304 return event.data.message.content
305 .filter((block): block is ContentBlock & { type: 'text' } => block.type === 'text')
306 .map(block => block.text)
307 .join('')
308 }
309 return ''
310}