server.ts 12 KB

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