server.ts 10 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261
  1. /**
  2. * JSON-RPC method and notification surface for out-of-process harness SDKs.
  3. * The surrounding context owns plugins, persistence, and configured adapters.
  4. *
  5. * @module @deepseek-ai/dsh-jsonrpc/server
  6. */
  7. import type { Context } from 'cordis'
  8. import { resolve } from 'node:path'
  9. import type { ContentBlock } from '@deepseek-ai/dsh-llm'
  10. import type { Agent, AgentHandle } from '@deepseek-ai/dsh-agent'
  11. import { carrierKeyOf, type Scoped } from '@deepseek-ai/dsh-scope'
  12. import { findLastMessageTurnEnd, SessionId, type TurnEndReason } from '@deepseek-ai/dsh-session'
  13. import type SubagentService from '@deepseek-ai/dsh-subagent'
  14. import type { SubagentRunEndInfo } from '@deepseek-ai/dsh-subagent'
  15. import * as LlmDeepSeek from '@deepseek-ai/dsh-llm-deepseek'
  16. import type { JsonRpcTransportPeer } from './transport.ts'
  17. /** Parameters for the process-wide SDK handshake. */
  18. export interface InitializeParams {
  19. /** Working directory recorded on every SDK-created session's header. */
  20. cwd: string
  21. /** Provider route every SDK-created agent runs on. */
  22. provider: string
  23. /** Model name every SDK-created agent runs on (see {@link HarnessSdkServer.initialize} for adapter fallback). */
  24. model: string
  25. }
  26. /** Wire-stable server identity returned by initialization. */
  27. export interface InitializeResult {
  28. /** Wire-stable server identity (`deepseek-harness-sdk-runtime`) and version. */
  29. serverInfo: { name: string; version: string }
  30. }
  31. /** One user turn on one SDK session. */
  32. export interface SessionPromptParams {
  33. /** The SDK-side session id; an unknown id lazily creates the agent+session pair. */
  34. sessionId: string
  35. /** The prompt content blocks, sent verbatim as the user message. */
  36. contentBlocks: ContentBlock[]
  37. }
  38. /** Prompt acceptance after turn settlement; outcome rides on `session.finished`. */
  39. export interface SessionPromptResult {
  40. /** Always `true`; the turn outcome is the paired `session.finished` notification. */
  41. accepted: true
  42. }
  43. interface SessionRecord {
  44. handle: AgentHandle
  45. lastTurnEnd: TurnEndReason | undefined
  46. activePrompt: boolean
  47. }
  48. /** Recover the delegating parent from the service-owned scoped carrier. */
  49. function subagentParentOf(carrier: Scoped<SubagentService>): Agent {
  50. return carrierKeyOf(carrier) as Agent
  51. }
  52. /** Deployment-specific status mapping for SDK turn and subagent outcomes. */
  53. export interface HarnessSdkServerOptions {
  54. /** Report max-token termination as an accepted result instead of an infrastructure error. */
  55. maxTokensAsSuccess?: boolean
  56. }
  57. function successStatus(reason: string, options: HarnessSdkServerOptions): '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 HarnessSdkServer {
  67. private cwd = process.cwd()
  68. private provider = 'deepseek'
  69. private model = 'deepseek'
  70. private llmFiber: { dispose(): Promise<void> } | undefined
  71. private readonly sessions = new Map<string, SessionRecord>()
  72. private readonly sessionCreations = new Map<string, Promise<SessionRecord>>()
  73. private readonly disposers: (() => void)[] = []
  74. private shutdownTask: Promise<Record<string, never>> | undefined
  75. private shuttingDown = false
  76. constructor(
  77. private readonly ctx: Context,
  78. private readonly transport: JsonRpcTransportPeer,
  79. private readonly options: HarnessSdkServerOptions = {},
  80. ) {
  81. const serverOptions = this.options
  82. this.disposers.push(ctx.on('session/event', (session, event) => {
  83. if (event.type === 'turn/end') {
  84. const rec = this.sessions.get(String(session.id))
  85. if (rec && findLastMessageTurnEnd(session.events)?.seq === event.seq) {
  86. rec.lastTurnEnd = event.data.reason
  87. }
  88. }
  89. this.transport.notify('session.event', { sessionId: String(session.id), event })
  90. }))
  91. this.disposers.push(ctx.on('session/created', (session) => {
  92. const parentSession = session.header.parentSession
  93. if (parentSession === undefined) return
  94. this.transport.notify('subagent.started', {
  95. parentSessionId: String(parentSession),
  96. childSessionId: String(session.id),
  97. })
  98. }))
  99. this.disposers.push(ctx.on('subagent/end', function (this: Scoped<SubagentService>, info: SubagentRunEndInfo) {
  100. const parent = subagentParentOf(this)
  101. // This protocol reports only in-process child sessions. The service
  102. // snapshots the provider's exact run provenance through child disposal;
  103. // matching ids or parent lineage alone never establishes locality.
  104. if (!info.local) return
  105. transport.notify('subagent.finished', {
  106. provider: info.provider,
  107. agentId: String(info.id),
  108. parentSessionId: String(parent.session.id),
  109. childSessionId: String(info.id),
  110. status: successStatus(info.stopReason, serverOptions),
  111. stopReason: info.stopReason,
  112. ...(info.lastAssistantMessage === undefined ? {} : { lastAssistantMessage: info.lastAssistantMessage }),
  113. })
  114. }))
  115. }
  116. /**
  117. * Configure the SDK route, mounting the DeepSeek fallback only when unowned.
  118. * @param params - SDK handshake parameters.
  119. * @returns server identity for the handshake.
  120. */
  121. async initialize(params: InitializeParams): Promise<InitializeResult> {
  122. this.cwd = resolve(params.cwd)
  123. this.provider = params.provider
  124. this.model = params.model
  125. if (!this.hasAdapterFor(this.provider)) {
  126. if (this.provider !== 'deepseek') throw new Error(`no adapter registered for provider "${this.provider}"`)
  127. this.llmFiber = await this.ctx.plugin(LlmDeepSeek, {})
  128. }
  129. return { serverInfo: { name: 'deepseek-harness-sdk-runtime', version: '0.0.1' } }
  130. }
  131. /**
  132. * Run one prompt to settlement; overlap on the same session fails.
  133. * @param params - target session and user content.
  134. * @returns acceptance after the turn settled.
  135. */
  136. async prompt(params: SessionPromptParams): Promise<SessionPromptResult> {
  137. const rec = await this.getOrCreateSession(params.sessionId)
  138. if (rec.activePrompt) throw new Error(`session already has an active prompt: ${params.sessionId}`)
  139. rec.activePrompt = true
  140. try {
  141. rec.lastTurnEnd = undefined
  142. rec.handle.agent.followup(params.contentBlocks)
  143. await rec.handle.agent.whenIdle()
  144. const status = this.finishedStatus(rec.lastTurnEnd)
  145. this.transport.notify('session.finished', {
  146. sessionId: params.sessionId,
  147. status,
  148. reason: rec.lastTurnEnd,
  149. })
  150. return { accepted: true }
  151. } finally {
  152. rec.activePrompt = false
  153. }
  154. }
  155. /**
  156. * Dispose server-owned agents, adapter, and subscriptions to quiescence.
  157. * The surrounding context remains running.
  158. * @returns empty JSON-RPC result.
  159. */
  160. shutdown(): Promise<Record<string, never>> {
  161. this.shutdownTask ??= this.performShutdown()
  162. return this.shutdownTask
  163. }
  164. private async performShutdown(): Promise<Record<string, never>> {
  165. this.shuttingDown = true
  166. const pendingCreations = [...this.sessionCreations.values()]
  167. await Promise.allSettled(pendingCreations)
  168. this.sessionCreations.clear()
  169. const records = [...this.sessions.values()]
  170. this.sessions.clear()
  171. const failures: unknown[] = []
  172. while (this.disposers.length > 0) {
  173. try {
  174. this.disposers.pop()?.()
  175. } catch (error) {
  176. failures.push(error)
  177. }
  178. }
  179. const teardownResults = await Promise.allSettled([
  180. ...records.map(rec => Promise.resolve().then(() => rec.handle.dispose())),
  181. ...(this.llmFiber === undefined ? [] : [Promise.resolve().then(() => this.llmFiber?.dispose())]),
  182. ])
  183. this.llmFiber = undefined
  184. failures.push(...teardownResults
  185. .filter((result): result is PromiseRejectedResult => result.status === 'rejected')
  186. .map(result => result.reason as unknown))
  187. if (failures.length === 1) throw failures[0]
  188. if (failures.length > 1) throw new AggregateError(failures, 'SDK server teardown failed')
  189. return {}
  190. }
  191. /**
  192. * Dispatch one incoming JSON-RPC request to its typed handler. Throws (→ a
  193. * JSON-RPC error response) on an unknown method.
  194. * @param method - the JSON-RPC method name.
  195. * @param params - the raw params object from the wire.
  196. * @returns the handler's result, to be serialized as the response.
  197. */
  198. async handleRequest(method: string, params: Record<string, unknown> | undefined): Promise<unknown> {
  199. switch (method) {
  200. case 'initialize':
  201. return this.initialize(params as unknown as InitializeParams)
  202. case 'session/prompt':
  203. return this.prompt(params as unknown as SessionPromptParams)
  204. case 'shutdown':
  205. return this.shutdown()
  206. default:
  207. throw new Error(`unknown DeepSeek Harness SDK runtime method: ${method}`)
  208. }
  209. }
  210. private async getOrCreateSession(sessionId: string): Promise<SessionRecord> {
  211. if (this.shuttingDown) throw new Error('SDK server is shutting down')
  212. const existing = this.sessions.get(sessionId)
  213. if (existing) return existing
  214. const pending = this.sessionCreations.get(sessionId)
  215. if (pending) return pending
  216. const creation = this.createSession(sessionId)
  217. this.sessionCreations.set(sessionId, creation)
  218. void creation.then(
  219. () => { this.sessionCreations.delete(sessionId) },
  220. () => { this.sessionCreations.delete(sessionId) },
  221. )
  222. return creation
  223. }
  224. private async createSession(sessionId: string): Promise<SessionRecord> {
  225. const handle = await this.ctx.agents.create({
  226. sessionId: SessionId(sessionId),
  227. meta: { cwd: this.cwd },
  228. agentOptions: { provider: this.provider, model: this.model },
  229. })
  230. const rec: SessionRecord = { handle, lastTurnEnd: undefined, activePrompt: false }
  231. this.sessions.set(sessionId, rec)
  232. return rec
  233. }
  234. private finishedStatus(reason: TurnEndReason | undefined): 'ok' | 'error' {
  235. if (!reason) return 'error'
  236. return successStatus(reason.kind, this.options)
  237. }
  238. private hasAdapterFor(provider: string): boolean {
  239. return this.ctx.get('llm')?.listProviders().some(entry => entry.id === provider) ?? false
  240. }
  241. }