返回源码地图

packages/sdk/server/src/server.ts

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

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

1/**
2 * JSON-RPC methods and notifications for out-of-process harness SDKs.
3 * The surrounding context owns plugins, persistence, and configured adapters.
4 *
5 * @module @deepseek-ai/dsh-sdk-jsonrpc-server/server
6 */
7
8import type { Context } from '@deepseek-ai/cordis'
9import { resolve } from 'node:path'
10import { brandString } from '@deepseek-ai/dsh-brand'
11import type { Agent, AgentHandle } from '@deepseek-ai/dsh-agent'
12import { admitEncodedImages, type EncodedImageAttachment, type ImageAttachmentRef } from '@deepseek-ai/dsh-attachment'
13import { createUserMessage, ReasoningEffortId, type ContentBlock, type LlmRuntime } from '@deepseek-ai/dsh-llm'
14import { carrierKeyOf, type Scoped } from '@deepseek-ai/dsh-scope'
15import type { SessionId } from '@deepseek-ai/dsh-session'
16import type SubagentRuntime from '@deepseek-ai/dsh-subagent'
17import type { SubagentRunEndInfo } from '@deepseek-ai/dsh-subagent'
18import * as LlmDeepSeek from '@deepseek-ai/dsh-llm-deepseek-api-key'
19import type {
20 InitializeParams,
21 InitializeResult,
22 JsonRpcTransportPeer,
23 SessionEventNotification,
24 SessionPromptParams,
25 SessionPromptResult,
26 SdkEncodedImageBlock,
27 SubagentFinishedNotification,
28 SubagentStartedNotification,
29} from '@deepseek-ai/dsh-sdk-protocol'
30
31interface SessionRecord {
32 handle: AgentHandle
33}
34
35function encodedImage(block: SessionPromptParams['contentBlocks'][number]): block is SdkEncodedImageBlock {
36 return block.type === 'image' && 'data' in block
37}
38
39async function durablePromptContent(ctx: Context, blocks: SessionPromptParams['contentBlocks']): Promise<ContentBlock[]> {
40 const images = blocks.filter(encodedImage)
41 if (images.length === 0) return blocks as ContentBlock[]
42 const attachments = ctx.get('attachments')
43 if (attachments === undefined) throw new Error('SDK image prompt requires an attachment store')
44 const refs = await admitEncodedImages(attachments, images.map((image): EncodedImageAttachment => ({
45 data: image.data,
46 mediaType: image.mimeType,
47 })))
48 let next = 0
49 return blocks.map(block => encodedImage(block)
50 ? { type: 'image', attachment: refs[next++] as ImageAttachmentRef }
51 : block)
52}
53
54/** Recover the delegating parent from the service-owned scoped carrier. */
55function subagentParentOf(carrier: Scoped<SubagentRuntime>): Agent {
56 return carrierKeyOf(carrier) as Agent
57}
58
59/** Deployment-specific status mapping for SDK turn and subagent outcomes. */
60export interface HarnessSdkJsonRpcServerOptions {
61 /** Report max-token termination as an accepted result instead of an infrastructure error. */
62 maxTokensAsSuccess?: boolean
63}
64
65function successStatus(reason: string, options: HarnessSdkJsonRpcServerOptions): 'ok' | 'error' {
66 if (reason === 'completed') return 'ok'
67 return reason === 'max-tokens' && options.maxTokensAsSuccess === true ? 'ok' : 'error'
68}
69
70/**
71 * SDK server over one booted harness context and transport peer. Construction
72 * subscribes to session, agent, and subagent lifecycle events until shutdown;
73 * reinitialization is unsupported.
74 */
75export class HarnessSdkJsonRpcServer {
76 private cwd = process.cwd()
77 private provider = 'deepseek-official'
78 private model = 'deepseek-official'
79 private reasoningEffort: ReturnType<typeof ReasoningEffortId> | undefined
80 private maxTokens: number | undefined
81 private llmFiber: { dispose(): Promise<void> } | undefined
82 private readonly sessions = new Map<string, SessionRecord>()
83 private readonly sessionCreations = new Map<string, Promise<SessionRecord>>()
84 private readonly disposers: (() => void)[] = []
85 private shutdownTask: Promise<Record<string, never>> | undefined
86 private shuttingDown = false
87 private initialized = false
88
89 constructor(
90 private readonly ctx: Context,
91 private readonly transport: JsonRpcTransportPeer,
92 private readonly options: HarnessSdkJsonRpcServerOptions = {},
93 ) {
94 const serverOptions = this.options
95 this.disposers.push(ctx.on('session/event', (session, event) => {
96 const payload: SessionEventNotification = { sessionId: String(session.id), event }
97 this.transport.notify('session.event', payload)
98 }))
99 this.disposers.push(ctx.on('agent/status', ({ agent, status }) => {
100 this.transport.notify('session.status', { sessionId: String(agent.session.id), status })
101 }))
102 this.disposers.push(ctx.on('session/created', (session) => {
103 const parentSession = session.header.parentSession
104 if (parentSession === undefined) return
105 const payload: SubagentStartedNotification = {
106 parentSessionId: String(parentSession),
107 childSessionId: String(session.id),
108 }
109 this.transport.notify('subagent.started', payload)
110 }))
111 this.disposers.push(ctx.on('subagent/end', function (this: Scoped<SubagentRuntime>, info: SubagentRunEndInfo) {
112 const parent = subagentParentOf(this)
113 // This protocol reports only in-process child sessions. The service
114 // snapshots the provider name and local flag through child disposal;
115 // matching ids or parent lineage alone never establishes locality.
116 if (!info.local) return
117 const payload: SubagentFinishedNotification = {
118 provider: info.provider,
119 agentId: String(info.id),
120 parentSessionId: String(parent.session.id),
121 childSessionId: String(info.id),
122 status: successStatus(info.stopReason, serverOptions),
123 stopReason: info.stopReason,
124 ...(info.lastAssistantMessage === undefined
125 ? {}
126 : { lastAssistantMessage: [...info.lastAssistantMessage] }),
127 }
128 transport.notify('subagent.finished', payload)
129 }))
130 }
131
132 /**
133 * Validate and configure the SDK route, mounting the DeepSeek fallback only when unowned.
134 * @param params - SDK handshake parameters.
135 * @returns server identity for the handshake.
136 */
137 async initialize(params: InitializeParams): Promise<InitializeResult> {
138 if (params.reasoningEffort !== undefined
139 && (typeof params.reasoningEffort !== 'string' || params.reasoningEffort.length === 0)) {
140 throw new TypeError('initialize reasoningEffort must be a non-empty string')
141 }
142 if (params.maxTokens !== undefined
143 && (!Number.isSafeInteger(params.maxTokens) || params.maxTokens <= 0)) {
144 throw new TypeError('initialize maxTokens must be a positive safe integer')
145 }
146 const cwd = resolve(params.cwd)
147 const provider = params.provider
148 const model = params.model
149 const reasoningEffort = params.reasoningEffort === undefined
150 ? undefined
151 : ReasoningEffortId(params.reasoningEffort)
152 if (!this.hasAdapterFor(provider)) {
153 if (provider !== 'deepseek-official') throw new Error(`no adapter registered for provider "${provider}"`)
154 this.llmFiber = await this.ctx.plugin(LlmDeepSeek)
155 }
156 // Adapter presence was read from this service above; a successful fallback mount also requires it.
157 const llm = this.ctx.get('llm') as LlmRuntime
158 await llm.resolveCallConfig({
159 provider,
160 model,
161 ...reasoningEffort === undefined ? {} : { reasoningEffort },
162 ...params.maxTokens === undefined ? {} : { maxTokens: params.maxTokens },
163 })
164 this.cwd = cwd
165 this.provider = provider
166 this.model = model
167 this.reasoningEffort = reasoningEffort
168 this.maxTokens = params.maxTokens
169 this.initialized = true
170 return { serverInfo: { name: 'deepseek-harness-sdk-runtime', version: '0.0.1' } }
171 }
172
173 /**
174 * Queue one identified prompt without assigning later activity to it.
175 * @param params - target session and user content.
176 * @returns the durable message identity.
177 */
178 async prompt(params: SessionPromptParams): Promise<SessionPromptResult> {
179 if (!this.initialized) throw new Error('SDK server is not initialized')
180 const rec = await this.getOrCreateSession(params.sessionId)
181 // An agent-loop-only reload disposes the loop's agents while this record
182 // survives; a retained agent accepts followup() silently, so validate the
183 // record against the live registry before delivery.
184 this.assertLiveAgent(rec, params.sessionId)
185 const content = await durablePromptContent(this.ctx, params.contentBlocks)
186 // Attachment admission crosses an async boundary where shutdown or an
187 // agent-loop reload may detach the retained handle.
188 this.assertLiveAgent(rec, params.sessionId)
189 const message = createUserMessage({
190 content,
191 source: { kind: 'user' },
192 })
193 rec.handle.agent.followup(message)
194 return { messageId: message.id }
195 }
196
197 private assertLiveAgent(rec: SessionRecord, sessionId: string): void {
198 if (this.ctx.agents.get(rec.handle.agent.id) !== rec.handle.agent) {
199 throw new Error(`session agent was disposed outside the server: ${sessionId}`)
200 }
201 }
202
203 /**
204 * Dispose server-owned agents, adapter, and subscriptions to quiescence.
205 * The surrounding context remains running.
206 * @returns empty JSON-RPC result.
207 */
208 shutdown(): Promise<Record<string, never>> {
209 this.shutdownTask ??= this.performShutdown()
210 return this.shutdownTask
211 }
212
213 private async performShutdown(): Promise<Record<string, never>> {
214 this.shuttingDown = true
215 const pendingCreations = [...this.sessionCreations.values()]
216 await Promise.allSettled(pendingCreations)
217 this.sessionCreations.clear()
218 const records = [...this.sessions.values()]
219 this.sessions.clear()
220 const failures: unknown[] = []
221 while (this.disposers.length > 0) {
222 try {
223 this.disposers.pop()?.()
224 } catch (error) {
225 failures.push(error)
226 }
227 }
228 const teardownResults = await Promise.allSettled([
229 ...records.map(rec => Promise.resolve().then(() => rec.handle.dispose())),
230 ...(this.llmFiber === undefined ? [] : [Promise.resolve().then(() => this.llmFiber?.dispose())]),
231 ])
232 this.llmFiber = undefined
233 failures.push(...teardownResults
234 .filter((result): result is PromiseRejectedResult => result.status === 'rejected')
235 .map(result => result.reason as unknown))
236 if (failures.length === 1) throw failures[0]
237 if (failures.length > 1) throw new AggregateError(failures, 'SDK server teardown failed')
238 return {}
239 }
240
241 /**
242 * Dispatch one incoming JSON-RPC request to its typed handler. Throws (→ a
243 * JSON-RPC error response) on an unknown method.
244 * @param method - the JSON-RPC method name.
245 * @param params - the raw params object from the wire.
246 * @returns the handler's result, to be serialized as the response.
247 */
248 async handleRequest(method: string, params: Record<string, unknown> | undefined): Promise<unknown> {
249 switch (method) {
250 case 'initialize':
251 return this.initialize(params as unknown as InitializeParams)
252 case 'session/prompt':
253 return this.prompt(params as unknown as SessionPromptParams)
254 case 'shutdown':
255 return this.shutdown()
256 default:
257 throw new Error(`unknown DeepSeek Harness SDK runtime method: ${method}`)
258 }
259 }
260
261 private async getOrCreateSession(sessionId: string): Promise<SessionRecord> {
262 if (this.shuttingDown) throw new Error('SDK server is shutting down')
263 const existing = this.sessions.get(sessionId)
264 if (existing) return existing
265 const pending = this.sessionCreations.get(sessionId)
266 if (pending) return pending
267 const creation = this.createSession(sessionId)
268 this.sessionCreations.set(sessionId, creation)
269 void creation.then(
270 () => { this.sessionCreations.delete(sessionId) },
271 () => { this.sessionCreations.delete(sessionId) },
272 )
273 return creation
274 }
275
276 private async createSession(sessionId: string): Promise<SessionRecord> {
277 // No preset composition: this server's compositions keep the model-facing
278 // rows in the host plane, so this agent reads them from the global layer. A
279 // deployment that configures a roster has to join one here first
280 // (@deepseek-ai/dsh-agent-preset-registry README, "Composing a child agent").
281 const handle = await this.ctx.agents.create({
282 sessionId: brandString<SessionId>(sessionId),
283 meta: { cwd: this.cwd },
284 agentOptions: {
285 provider: this.provider,
286 model: this.model,
287 ...this.reasoningEffort === undefined ? {} : { reasoningEffort: this.reasoningEffort },
288 ...this.maxTokens === undefined ? {} : { maxTokens: this.maxTokens },
289 },
290 })
291 const rec: SessionRecord = { handle }
292 this.sessions.set(sessionId, rec)
293 return rec
294 }
295
296 private hasAdapterFor(provider: string): boolean {
297 return this.ctx.get('llm')?.listProviders().some(entry => entry.id === provider) ?? false
298 }
299}