server.spec.ts 49 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203
  1. import { createUserMessage, LlmAdapter, ReasoningEffortId } from '@deepseek-ai/dsh-llm'
  2. import type { GenerateOptions, LlmResolvedModelInfo, StreamChunk } from '@deepseek-ai/dsh-llm'
  3. import { createServer } from 'node:http'
  4. import type { IncomingMessage, Server, ServerResponse } from 'node:http'
  5. import { mkdtemp, rm } from 'node:fs/promises'
  6. import { join } from 'node:path'
  7. import { tmpdir } from 'node:os'
  8. import { afterEach, describe, expect, it, vi } from 'vitest'
  9. import { Context } from '@deepseek-ai/cordis'
  10. import AgentRegistry, { type Agent, type AgentHandle } from '@deepseek-ai/dsh-agent'
  11. import AgentLoop from '@deepseek-ai/dsh-agent-loop'
  12. import { mountAgentLoopTestDependencies } from '@deepseek-ai/dsh-agent-loop-testkit'
  13. import SessionStore, { SessionId } from '@deepseek-ai/dsh-session'
  14. import JsonlSessionPersistence from '@deepseek-ai/dsh-session-persistence-jsonl'
  15. import * as LlmDeepSeek from '@deepseek-ai/dsh-llm-deepseek'
  16. import SubagentRuntime, { type SubagentResult, type SubagentRunEndInfo } from '@deepseek-ai/dsh-subagent'
  17. import type { JsonRpcTransportPeer } from '@deepseek-ai/dsh-sdk-protocol'
  18. import { HarnessSdkJsonRpcServer } from '../src/index.ts'
  19. class FakeTransport implements JsonRpcTransportPeer {
  20. notifications: { method: string; params?: Record<string, unknown> }[] = []
  21. async request(method: string, params: object): Promise<unknown> {
  22. throw new Error(`the SDK server should not call host JSON-RPC method ${method} with ${JSON.stringify(params)}`)
  23. }
  24. notify(method: string, params?: object): void {
  25. this.notifications.push(params === undefined ? { method } : { method, params: params as Record<string, unknown> })
  26. }
  27. }
  28. const servers: Server[] = []
  29. afterEach(async () => {
  30. await Promise.all(servers.splice(0).map(server => new Promise(resolve => server.close(resolve))))
  31. vi.unstubAllEnvs()
  32. })
  33. async function mockCompletionServer(): Promise<{ url: string; requests: unknown[]; headers: IncomingMessage['headers'][] }> {
  34. const requests: unknown[] = []
  35. const headers: IncomingMessage['headers'][] = []
  36. const server = createServer((request: IncomingMessage, response: ServerResponse) => {
  37. let body = ''
  38. request.on('data', (chunk: Buffer) => { body += chunk.toString('utf8') })
  39. request.on('end', () => {
  40. requests.push(JSON.parse(body))
  41. headers.push(request.headers)
  42. response.writeHead(200, { 'content-type': 'text/event-stream' })
  43. response.write('data: {"choices":[{"delta":{"role":"assistant","content":null,"reasoning_content":""}}]}\n\n')
  44. response.write('data: {"choices":[{"delta":{"content":"done"}}]}\n\n')
  45. response.write('data: {"choices":[{"delta":{"content":""},"finish_reason":"stop"}],"usage":{"prompt_tokens":3,"completion_tokens":1}}\n\n')
  46. response.write('data: [DONE]\n\n')
  47. response.end()
  48. })
  49. })
  50. servers.push(server)
  51. await new Promise<void>(resolve => server.listen(0, '127.0.0.1', resolve))
  52. const address = server.address()
  53. if (address === null || typeof address === 'string') throw new Error('no port')
  54. return { url: `http://127.0.0.1:${address.port}`, requests, headers }
  55. }
  56. async function makeHarness(storageDir: string) {
  57. const ctx = new Context()
  58. await mountAgentLoopTestDependencies(ctx)
  59. await ctx.plugin(AgentLoop, { agents: [] })
  60. await ctx.plugin(SubagentRuntime)
  61. await ctx.plugin(JsonlSessionPersistence, { root: storageDir })
  62. await new Promise(resolve => setTimeout(resolve, 50))
  63. return ctx
  64. }
  65. /** Drive the owning service so test lifecycle events carry the real parent scope. */
  66. async function settleSubagent(
  67. ctx: Context,
  68. parent: Agent,
  69. info: Omit<SubagentRunEndInfo, 'runId' | 'local'> & { localAgent: Agent | undefined },
  70. beforeSettle?: () => Promise<void>,
  71. ): Promise<void> {
  72. const result = Promise.withResolvers<SubagentResult>()
  73. const disposeProvider = ctx.subagents.registerProvider({
  74. name: info.provider,
  75. capabilities: { agentOptions: false, outputSchema: false, depthLimit: false, toolFilter: false, persona: false },
  76. inheritsParentContext: false,
  77. async start() {
  78. return {
  79. id: info.id,
  80. localAgent: info.localAgent,
  81. result: result.promise,
  82. dispose: () => Promise.resolve(),
  83. }
  84. },
  85. })
  86. try {
  87. const run = await ctx.subagents.start(info.provider, {
  88. parent,
  89. prompt: [],
  90. signal: new AbortController().signal,
  91. })
  92. await beforeSettle?.()
  93. if (info.lastAssistantMessage === undefined) {
  94. result.reject(new Error('synthetic infrastructure failure'))
  95. } else {
  96. result.resolve({ output: info.lastAssistantMessage, stopReason: info.stopReason })
  97. }
  98. await run.result.then(() => undefined, () => undefined)
  99. await run.dispose()
  100. } finally {
  101. disposeProvider()
  102. }
  103. }
  104. describe('HarnessSdkJsonRpcServer', () => {
  105. it('creates a harness agent and calls the configured OpenAI-compatible endpoint', { timeout: 15_000 }, async () => {
  106. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-'))
  107. const llmServer = await mockCompletionServer()
  108. vi.stubEnv('DEEPSEEK_API_KEY', 'test-key')
  109. vi.stubEnv('DEEPSEEK_BASE_URL', llmServer.url)
  110. const ctx = await makeHarness(storageDir)
  111. try {
  112. const transport = new FakeTransport()
  113. const server = new HarnessSdkJsonRpcServer(ctx, transport)
  114. const init = await server.handleRequest('initialize', {
  115. cwd: storageDir,
  116. provider: 'deepseek-official',
  117. model: 'dsagent-model',
  118. reasoningEffort: 'max',
  119. maxTokens: 321,
  120. }) as { serverInfo: { name: string } }
  121. expect(init.serverInfo.name).toBe('deepseek-harness-sdk-runtime')
  122. const receipt = await server.handleRequest('session/prompt', {
  123. sessionId: 'main',
  124. contentBlocks: [{ type: 'text', text: 'fix it' }],
  125. })
  126. expect((receipt as { messageId?: unknown }).messageId).toBeTypeOf('string')
  127. await vi.waitFor(() => { expect(llmServer.requests).toHaveLength(1) })
  128. const body = llmServer.requests[0] as {
  129. model: string
  130. messages: { role: string }[]
  131. reasoning_effort?: string
  132. max_tokens?: number
  133. }
  134. expect(body.model).toBe('dsagent-model')
  135. expect(body.reasoning_effort).toBe('max')
  136. expect(body.max_tokens).toBe(321)
  137. expect(body.messages[0]?.role).toBe('system')
  138. expect(body.messages.at(-1)?.role).toBe('user')
  139. expect(llmServer.headers[0]?.authorization).toBe('Bearer test-key')
  140. expect(transport.notifications.some(n => n.method === 'session.event')).toBe(true)
  141. await vi.waitFor(() => {
  142. expect(transport.notifications.findLast(n => n.method === 'session.status')).toEqual({
  143. method: 'session.status',
  144. params: { sessionId: 'main', status: 'idle' },
  145. })
  146. })
  147. await server.handleRequest('session/prompt', {
  148. sessionId: 'main',
  149. contentBlocks: [{ type: 'text', text: 'again' }],
  150. })
  151. await vi.waitFor(() => { expect(llmServer.requests).toHaveLength(2) })
  152. const orphanHandle = await ctx.agents.create({
  153. sessionId: SessionId('orphan-session'),
  154. meta: { cwd: storageDir },
  155. agentOptions: { provider: 'deepseek-official', model: 'dsagent-model' },
  156. })
  157. orphanHandle.agent.followup(createUserMessage({ content: [{ type: 'text', text: 'outside the sdk session map' }], source: { kind: 'user' } }))
  158. await orphanHandle.agent.whenIdle()
  159. await orphanHandle.dispose()
  160. expect(llmServer.requests).toHaveLength(3)
  161. await server.handleRequest('shutdown', undefined)
  162. } finally {
  163. await ctx.fiber.dispose()
  164. await rm(storageDir, { recursive: true, force: true })
  165. }
  166. })
  167. it('queues overlapping prompts for one session without blocking other sessions', async () => {
  168. const mainFollowup = vi.fn<Agent['followup']>()
  169. const mainAgent = ({
  170. id: SessionId('main'),
  171. followup: mainFollowup,
  172. } satisfies Pick<Agent, 'id' | 'followup'>) as unknown as Agent
  173. const otherFollowup = vi.fn<Agent['followup']>()
  174. const otherAgent = ({
  175. id: SessionId('other'),
  176. followup: otherFollowup,
  177. } satisfies Pick<Agent, 'id' | 'followup'>) as unknown as Agent
  178. const mainHandle = { agent: mainAgent, dispose: vi.fn(() => Promise.resolve()) }
  179. const otherHandle = { agent: otherAgent, dispose: vi.fn(() => Promise.resolve()) }
  180. const create = vi.fn(async (options: { sessionId: SessionId }) =>
  181. String(options.sessionId) === 'main' ? mainHandle : otherHandle)
  182. const liveAgents = new Map<string, Agent>([['main', mainAgent], ['other', otherAgent]])
  183. const ctx = {
  184. on: vi.fn(() => () => undefined),
  185. agents: { create, get: (id: SessionId) => liveAgents.get(String(id)) },
  186. get: () => undefined,
  187. } as unknown as Context
  188. const server = new HarnessSdkJsonRpcServer(ctx, new FakeTransport())
  189. // This isolated prompt test begins after the handshake boundary.
  190. ;(server as unknown as { initialized: boolean }).initialized = true
  191. const prompt = (sessionId: string, text: string) => server.prompt({
  192. sessionId,
  193. contentBlocks: [{ type: 'text', text }],
  194. })
  195. expect((await prompt('main', 'first')).messageId).toBeTypeOf('string')
  196. expect((await prompt('main', 'overlap')).messageId).toBeTypeOf('string')
  197. expect((await prompt('other', 'independent')).messageId).toBeTypeOf('string')
  198. expect(mainFollowup).toHaveBeenCalledTimes(2)
  199. expect(otherFollowup).toHaveBeenCalledOnce()
  200. await server.shutdown()
  201. expect(mainHandle.dispose).toHaveBeenCalledOnce()
  202. expect(otherHandle.dispose).toHaveBeenCalledOnce()
  203. })
  204. it('admits inline SDK images before the user message enters the session', async () => {
  205. const followup = vi.fn<Agent['followup']>()
  206. const agent = ({ id: SessionId('image'), followup } satisfies Pick<Agent, 'id' | 'followup'>) as unknown as Agent
  207. const handle = { agent, dispose: vi.fn(() => Promise.resolve()) }
  208. const ref = {
  209. attachmentId: 'sha256:image',
  210. mediaType: 'image/png',
  211. bytes: 1,
  212. width: 1,
  213. height: 1,
  214. }
  215. const saveImages = vi.fn(async () => [ref])
  216. const ctx = {
  217. on: vi.fn(() => () => undefined),
  218. agents: { create: vi.fn(async () => handle), get: () => agent },
  219. get: (name: string) => name === 'attachments' ? { saveImages } : undefined,
  220. } as unknown as Context
  221. const server = new HarnessSdkJsonRpcServer(ctx, new FakeTransport())
  222. // This isolated prompt test begins after the handshake boundary.
  223. ;(server as unknown as { initialized: boolean }).initialized = true
  224. await server.prompt({
  225. sessionId: 'image',
  226. contentBlocks: [
  227. { type: 'text', text: 'inspect' },
  228. { type: 'image', data: 'AQ==', mimeType: 'image/png' },
  229. ],
  230. })
  231. expect(saveImages).toHaveBeenCalledWith([{ data: Uint8Array.of(1), mediaType: 'image/png' }])
  232. expect(followup.mock.calls[0]?.[0].content).toEqual([
  233. { type: 'text', text: 'inspect' },
  234. { type: 'image', attachment: ref },
  235. ])
  236. await server.shutdown()
  237. })
  238. it('rejects inline SDK images when the composition has no attachment store', async () => {
  239. const followup = vi.fn<Agent['followup']>()
  240. const agent = ({ id: SessionId('image'), followup } satisfies Pick<Agent, 'id' | 'followup'>) as unknown as Agent
  241. const handle = { agent, dispose: vi.fn(() => Promise.resolve()) }
  242. const ctx = {
  243. on: vi.fn(() => () => undefined),
  244. agents: { create: vi.fn(async () => handle), get: () => agent },
  245. get: () => undefined,
  246. } as unknown as Context
  247. const server = new HarnessSdkJsonRpcServer(ctx, new FakeTransport())
  248. // This isolated prompt test begins after the handshake boundary.
  249. ;(server as unknown as { initialized: boolean }).initialized = true
  250. await expect(server.prompt({
  251. sessionId: 'image',
  252. contentBlocks: [{ type: 'image', data: 'AQ==', mimeType: 'image/png' }],
  253. })).rejects.toThrow('SDK image prompt requires an attachment store')
  254. expect(followup).not.toHaveBeenCalled()
  255. await server.shutdown()
  256. })
  257. it('rechecks agent liveness after asynchronous image admission', async () => {
  258. const followup = vi.fn<Agent['followup']>()
  259. const agent = ({ id: SessionId('image-race'), followup } satisfies Pick<Agent, 'id' | 'followup'>) as unknown as Agent
  260. const handle = { agent, dispose: vi.fn(() => Promise.resolve()) }
  261. const admitted = Promise.withResolvers<Array<{
  262. attachmentId: string
  263. mediaType: string
  264. bytes: number
  265. }>>()
  266. const saveImages = vi.fn(() => admitted.promise)
  267. let live = true
  268. const ctx = {
  269. on: vi.fn(() => () => undefined),
  270. agents: {
  271. create: vi.fn(async () => handle),
  272. get: () => live ? agent : undefined,
  273. },
  274. get: (name: string) => name === 'attachments' ? { saveImages } : undefined,
  275. } as unknown as Context
  276. const server = new HarnessSdkJsonRpcServer(ctx, new FakeTransport())
  277. // This isolated prompt test begins after the handshake boundary.
  278. ;(server as unknown as { initialized: boolean }).initialized = true
  279. const prompting = server.prompt({
  280. sessionId: 'image-race',
  281. contentBlocks: [{ type: 'image', data: 'AQ==', mimeType: 'image/png' }],
  282. })
  283. await vi.waitFor(() => { expect(saveImages).toHaveBeenCalledOnce() })
  284. live = false
  285. admitted.resolve([{ attachmentId: 'sha256:image', mediaType: 'image/png', bytes: 1 }])
  286. await expect(prompting).rejects.toThrow('session agent was disposed outside the server: image-race')
  287. expect(followup).not.toHaveBeenCalled()
  288. await server.shutdown()
  289. })
  290. it('rejects a prompt for a session whose agent was disposed outside the server', async () => {
  291. const followup = vi.fn<Agent['followup']>()
  292. const agent = ({
  293. id: SessionId('zombie'),
  294. followup,
  295. whenIdle: vi.fn(() => Promise.resolve()),
  296. } satisfies Pick<Agent, 'id' | 'followup' | 'whenIdle'>) as unknown as Agent
  297. const handle = { agent, dispose: vi.fn(() => Promise.resolve()) }
  298. // The registry drops the agent after creation, modelling an agent-loop-only
  299. // reload that leaves the server's SessionRecord pointing at a detached agent.
  300. let live = true
  301. const ctx = {
  302. on: vi.fn(() => () => undefined),
  303. agents: {
  304. create: vi.fn(async () => handle),
  305. get: (id: SessionId) => (live && String(id) === 'zombie' ? agent : undefined),
  306. },
  307. get: () => undefined,
  308. } as unknown as Context
  309. const server = new HarnessSdkJsonRpcServer(ctx, new FakeTransport())
  310. // This isolated prompt test begins after the handshake boundary.
  311. ;(server as unknown as { initialized: boolean }).initialized = true
  312. const prompt = (text: string) => server.prompt({
  313. sessionId: 'zombie',
  314. contentBlocks: [{ type: 'text', text }],
  315. })
  316. expect((await prompt('while live')).messageId).toBeTypeOf('string')
  317. live = false
  318. await expect(prompt('after detach')).rejects.toThrow('session agent was disposed outside the server: zombie')
  319. // The detached agent was never driven by the rejected prompt.
  320. expect(followup).toHaveBeenCalledOnce()
  321. await server.shutdown()
  322. })
  323. it('forwards whole-agent status without attributing a turn outcome', async () => {
  324. const ctx = new Context()
  325. await ctx.plugin(SessionStore)
  326. await ctx.plugin(AgentRegistry)
  327. const transport = new FakeTransport()
  328. const server = new HarnessSdkJsonRpcServer(ctx, transport)
  329. const session = ctx.sessions.create(SessionId('message-outcome'))
  330. const agent = ({
  331. id: SessionId('message-outcome'),
  332. session,
  333. } satisfies Pick<Agent, 'id' | 'session'>) as Agent
  334. ctx.emit('agent/status', { agent, status: 'running' })
  335. ctx.emit('agent/status', { agent, status: 'idle' })
  336. expect(transport.notifications.filter(notification => notification.method === 'session.status'))
  337. .toEqual([
  338. { method: 'session.status', params: { sessionId: 'message-outcome', status: 'running' } },
  339. { method: 'session.status', params: { sessionId: 'message-outcome', status: 'idle' } },
  340. ])
  341. await server.shutdown()
  342. await ctx.fiber.dispose()
  343. })
  344. it('notifies the host when a child session is created with parent lineage', async () => {
  345. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-subagent-'))
  346. const ctx = await makeHarness(storageDir)
  347. try {
  348. const transport = new FakeTransport()
  349. const server = new HarnessSdkJsonRpcServer(ctx, transport)
  350. ctx.sessions.create(SessionId('root-session'), {
  351. meta: { cwd: storageDir },
  352. })
  353. ctx.sessions.create(SessionId('child-session'), {
  354. meta: { cwd: storageDir, parentSession: SessionId('main') },
  355. })
  356. expect(transport.notifications).toContainEqual({
  357. method: 'subagent.started',
  358. params: {
  359. parentSessionId: 'main',
  360. childSessionId: 'child-session',
  361. },
  362. })
  363. await server.shutdown()
  364. } finally {
  365. await ctx.fiber.dispose()
  366. await rm(storageDir, { recursive: true, force: true })
  367. }
  368. })
  369. it('creates an SDK session without an optional system prompt', { timeout: 15_000 }, async () => {
  370. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-no-system-'))
  371. const llmServer = await mockCompletionServer()
  372. vi.stubEnv('DEEPSEEK_API_KEY', 'test-key')
  373. vi.stubEnv('DEEPSEEK_BASE_URL', llmServer.url)
  374. const ctx = await makeHarness(storageDir)
  375. try {
  376. const server = new HarnessSdkJsonRpcServer(ctx, new FakeTransport())
  377. await server.initialize({ cwd: storageDir, provider: 'deepseek-official', model: 'plain-model' })
  378. await server.prompt({
  379. sessionId: 'plain',
  380. contentBlocks: [{ type: 'text', text: 'hello' }],
  381. })
  382. await vi.waitFor(() => { expect(llmServer.requests).toHaveLength(1) })
  383. await server.shutdown()
  384. } finally {
  385. await ctx.fiber.dispose()
  386. await rm(storageDir, { recursive: true, force: true })
  387. }
  388. })
  389. it('notifies the host when a subagent run settles', async () => {
  390. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-subagent-end-'))
  391. const ctx = await makeHarness(storageDir)
  392. try {
  393. const transport = new FakeTransport()
  394. const server = new HarnessSdkJsonRpcServer(ctx, transport)
  395. const parentHandle = await ctx.agents.create({
  396. sessionId: SessionId('main'),
  397. meta: { cwd: storageDir },
  398. agentOptions: { provider: 'deepseek-official', model: 'deepseek-official' },
  399. })
  400. // A custom in-process provider may own its child at the provider/root
  401. // scope while preserving durable parent lineage.
  402. const handle = await ctx.agents.create({
  403. sessionId: SessionId('child-session'),
  404. meta: { cwd: storageDir, parentSession: SessionId('main') },
  405. agentOptions: { provider: 'deepseek-official', model: 'deepseek-official' },
  406. })
  407. expect(ctx.agents.roots()).toContain(handle.agent)
  408. const parentlessHandle = await parentHandle.agent.ctx.agents.create({
  409. sessionId: SessionId('parentless-child-session'),
  410. meta: { cwd: storageDir },
  411. agentOptions: { model: 'deepseek-official' },
  412. parentAgent: parentHandle.agent,
  413. })
  414. await settleSubagent(ctx, parentHandle.agent, {
  415. provider: 'spawn',
  416. id: SessionId('child-session'),
  417. localAgent: handle.agent,
  418. stopReason: 'completed',
  419. lastAssistantMessage: [{ type: 'text', text: 'child done' }],
  420. }, () => handle.dispose())
  421. await settleSubagent(ctx, parentHandle.agent, {
  422. provider: 'spawn',
  423. id: SessionId('parentless-child-session'),
  424. localAgent: parentlessHandle.agent,
  425. stopReason: 'error',
  426. }, () => parentlessHandle.dispose())
  427. expect(transport.notifications).toContainEqual({
  428. method: 'subagent.finished',
  429. params: {
  430. provider: 'spawn',
  431. agentId: 'child-session',
  432. parentSessionId: 'main',
  433. childSessionId: 'child-session',
  434. status: 'ok',
  435. stopReason: 'completed',
  436. lastAssistantMessage: [{ type: 'text', text: 'child done' }],
  437. },
  438. })
  439. expect(transport.notifications).toContainEqual({
  440. method: 'subagent.finished',
  441. params: {
  442. provider: 'spawn',
  443. agentId: 'parentless-child-session',
  444. parentSessionId: 'main',
  445. childSessionId: 'parentless-child-session',
  446. status: 'error',
  447. stopReason: 'error',
  448. },
  449. })
  450. await parentHandle.dispose()
  451. await server.shutdown()
  452. } finally {
  453. await ctx.fiber.dispose()
  454. await rm(storageDir, { recursive: true, force: true })
  455. }
  456. })
  457. it('ignores a remote run id that collides with a local child of the same parent', async () => {
  458. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-subagent-remote-collision-'))
  459. const ctx = await makeHarness(storageDir)
  460. try {
  461. const transport = new FakeTransport()
  462. const server = new HarnessSdkJsonRpcServer(ctx, transport)
  463. const parentHandle = await ctx.agents.create({
  464. sessionId: SessionId('collision-parent'),
  465. meta: { cwd: storageDir },
  466. agentOptions: { model: 'deepseek-official' },
  467. })
  468. const collidingChild = await parentHandle.agent.ctx.agents.create({
  469. sessionId: SessionId('remote-run-id'),
  470. meta: { cwd: storageDir, parentSession: SessionId('collision-parent') },
  471. agentOptions: { model: 'deepseek-official' },
  472. parentAgent: parentHandle.agent,
  473. })
  474. await settleSubagent(ctx, parentHandle.agent, {
  475. provider: 'remote',
  476. id: SessionId('remote-run-id'),
  477. localAgent: undefined,
  478. stopReason: 'completed',
  479. lastAssistantMessage: [],
  480. })
  481. expect(transport.notifications.some(notification =>
  482. notification.method === 'subagent.finished'
  483. && notification.params?.agentId === 'remote-run-id',
  484. )).toBe(false)
  485. await collidingChild.dispose()
  486. await parentHandle.dispose()
  487. await server.shutdown()
  488. } finally {
  489. await ctx.fiber.dispose()
  490. await rm(storageDir, { recursive: true, force: true })
  491. }
  492. })
  493. it('retains locality across continuation runs on one live child', async () => {
  494. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-subagent-continuation-'))
  495. const ctx = await makeHarness(storageDir)
  496. try {
  497. const transport = new FakeTransport()
  498. const server = new HarnessSdkJsonRpcServer(ctx, transport)
  499. const parentHandle = await ctx.agents.create({
  500. sessionId: SessionId('continuation-parent'),
  501. meta: { cwd: storageDir },
  502. agentOptions: { model: 'deepseek-official' },
  503. })
  504. const childHandle = await parentHandle.agent.ctx.agents.create({
  505. sessionId: SessionId('continuation-child'),
  506. meta: { cwd: storageDir, parentSession: SessionId('continuation-parent') },
  507. agentOptions: { model: 'deepseek-official' },
  508. parentAgent: parentHandle.agent,
  509. })
  510. await settleSubagent(ctx, parentHandle.agent, {
  511. provider: 'continuation',
  512. id: SessionId('continuation-child'),
  513. localAgent: childHandle.agent,
  514. stopReason: 'completed',
  515. lastAssistantMessage: [{ type: 'text', text: 'first' }],
  516. })
  517. await settleSubagent(ctx, parentHandle.agent, {
  518. provider: 'continuation',
  519. id: SessionId('continuation-child'),
  520. localAgent: childHandle.agent,
  521. stopReason: 'completed',
  522. lastAssistantMessage: [{ type: 'text', text: 'second' }],
  523. }, () => childHandle.dispose())
  524. expect(transport.notifications.filter(notification =>
  525. notification.method === 'subagent.finished'
  526. && notification.params?.childSessionId === 'continuation-child',
  527. )).toHaveLength(2)
  528. await parentHandle.dispose()
  529. await server.shutdown()
  530. } finally {
  531. await ctx.fiber.dispose()
  532. await rm(storageDir, { recursive: true, force: true })
  533. }
  534. })
  535. it('correlates reused local ids by parent scope when runs settle out of order', async () => {
  536. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-subagent-reuse-'))
  537. const ctx = await makeHarness(storageDir)
  538. try {
  539. const transport = new FakeTransport()
  540. const server = new HarnessSdkJsonRpcServer(ctx, transport)
  541. const oldParent = await ctx.agents.create({
  542. sessionId: SessionId('old-parent'),
  543. meta: { cwd: storageDir },
  544. agentOptions: { model: 'deepseek-official' },
  545. })
  546. const oldChild = await oldParent.agent.ctx.agents.create({
  547. sessionId: SessionId('reused-child'),
  548. meta: { cwd: storageDir, parentSession: SessionId('old-parent') },
  549. agentOptions: { model: 'deepseek-official' },
  550. parentAgent: oldParent.agent,
  551. })
  552. const first = Promise.withResolvers<SubagentResult>()
  553. const sameLifetime = Promise.withResolvers<SubagentResult>()
  554. const replacement = Promise.withResolvers<SubagentResult>()
  555. const results = [first.promise, sameLifetime.promise, replacement.promise]
  556. let starts = 0
  557. let currentLocalAgent = oldChild.agent
  558. const disposeProvider = ctx.subagents.registerProvider({
  559. name: 'reused',
  560. capabilities: { agentOptions: false, outputSchema: false, depthLimit: false, toolFilter: false, persona: false },
  561. inheritsParentContext: false,
  562. start() {
  563. const result = results[starts]
  564. starts += 1
  565. if (result === undefined) throw new Error('unexpected fourth reused-id run')
  566. return Promise.resolve({ id: SessionId('reused-child'), localAgent: currentLocalAgent, result, dispose: () => Promise.resolve() })
  567. },
  568. })
  569. const firstRun = await ctx.subagents.start('reused', {
  570. parent: oldParent.agent,
  571. prompt: [],
  572. signal: new AbortController().signal,
  573. })
  574. const sameLifetimeRun = await ctx.subagents.start('reused', {
  575. parent: oldParent.agent,
  576. prompt: [],
  577. signal: new AbortController().signal,
  578. })
  579. sameLifetime.resolve({ output: [{ type: 'text', text: 'same lifetime' }], stopReason: 'completed' })
  580. await sameLifetimeRun.result
  581. await oldChild.dispose()
  582. const newParent = await ctx.agents.create({
  583. sessionId: SessionId('new-parent'),
  584. meta: { cwd: storageDir },
  585. agentOptions: { model: 'deepseek-official' },
  586. })
  587. const newChild = await newParent.agent.ctx.agents.create({
  588. sessionId: SessionId('reused-child'),
  589. meta: { cwd: storageDir, parentSession: SessionId('new-parent') },
  590. agentOptions: { model: 'deepseek-official' },
  591. parentAgent: newParent.agent,
  592. })
  593. currentLocalAgent = newChild.agent
  594. const secondRun = await ctx.subagents.start('reused', {
  595. parent: newParent.agent,
  596. prompt: [],
  597. signal: new AbortController().signal,
  598. })
  599. replacement.resolve({ output: [{ type: 'text', text: 'new lifetime' }], stopReason: 'completed' })
  600. await secondRun.result
  601. first.resolve({ output: [{ type: 'text', text: 'old lifetime' }], stopReason: 'completed' })
  602. await firstRun.result
  603. await Promise.resolve()
  604. const finished = transport.notifications.filter(notification =>
  605. notification.method === 'subagent.finished'
  606. && notification.params?.childSessionId === 'reused-child',
  607. )
  608. expect(finished.map(notification => notification.params?.lastAssistantMessage)).toEqual([
  609. [{ type: 'text', text: 'same lifetime' }],
  610. [{ type: 'text', text: 'new lifetime' }],
  611. [{ type: 'text', text: 'old lifetime' }],
  612. ])
  613. expect(finished.map(notification => notification.params?.parentSessionId)).toEqual([
  614. 'old-parent',
  615. 'new-parent',
  616. 'old-parent',
  617. ])
  618. await firstRun.dispose()
  619. await sameLifetimeRun.dispose()
  620. await secondRun.dispose()
  621. disposeProvider()
  622. await newChild.dispose()
  623. await oldParent.dispose()
  624. await newParent.dispose()
  625. await server.shutdown()
  626. } finally {
  627. await ctx.fiber.dispose()
  628. await rm(storageDir, { recursive: true, force: true })
  629. }
  630. })
  631. it('keeps locality bound to the accepted run across provider re-registration', async () => {
  632. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-subagent-provider-reuse-'))
  633. const ctx = await makeHarness(storageDir)
  634. try {
  635. const transport = new FakeTransport()
  636. const server = new HarnessSdkJsonRpcServer(ctx, transport)
  637. const parent = await ctx.agents.create({
  638. sessionId: SessionId('provider-reuse-parent'),
  639. meta: { cwd: storageDir },
  640. agentOptions: { model: 'deepseek-official' },
  641. })
  642. const child = await parent.agent.ctx.agents.create({
  643. sessionId: SessionId('provider-reuse-child'),
  644. meta: { cwd: storageDir, parentSession: SessionId('provider-reuse-parent') },
  645. agentOptions: { model: 'deepseek-official' },
  646. parentAgent: parent.agent,
  647. })
  648. const localResult = Promise.withResolvers<SubagentResult>()
  649. const remoteResult = Promise.withResolvers<SubagentResult>()
  650. const unregisterLocal = ctx.subagents.registerProvider({
  651. name: 'reused-provider',
  652. capabilities: { agentOptions: false, outputSchema: false, depthLimit: false, toolFilter: false, persona: false },
  653. inheritsParentContext: false,
  654. start: () => Promise.resolve({
  655. id: SessionId('provider-reuse-child'),
  656. localAgent: child.agent,
  657. result: localResult.promise,
  658. dispose: () => Promise.resolve(),
  659. }),
  660. })
  661. const localRun = await ctx.subagents.start('reused-provider', {
  662. parent: parent.agent,
  663. prompt: [],
  664. signal: new AbortController().signal,
  665. })
  666. unregisterLocal()
  667. const unregisterRemote = ctx.subagents.registerProvider({
  668. name: 'reused-provider',
  669. capabilities: { agentOptions: false, outputSchema: false, depthLimit: false, toolFilter: false, persona: false },
  670. inheritsParentContext: false,
  671. start: () => Promise.resolve({
  672. id: SessionId('provider-reuse-child'),
  673. localAgent: undefined,
  674. result: remoteResult.promise,
  675. dispose: () => Promise.resolve(),
  676. }),
  677. })
  678. const remoteRun = await ctx.subagents.start('reused-provider', {
  679. parent: parent.agent,
  680. prompt: [],
  681. signal: new AbortController().signal,
  682. })
  683. remoteResult.resolve({ output: [{ type: 'text', text: 'remote' }], stopReason: 'completed' })
  684. await remoteRun.result
  685. await Promise.resolve()
  686. expect(transport.notifications.some(notification =>
  687. notification.method === 'subagent.finished'
  688. && notification.params?.lastAssistantMessage !== undefined,
  689. )).toBe(false)
  690. await child.dispose()
  691. localResult.resolve({ output: [{ type: 'text', text: 'local' }], stopReason: 'completed' })
  692. await localRun.result
  693. await Promise.resolve()
  694. expect(transport.notifications.filter(notification =>
  695. notification.method === 'subagent.finished'
  696. && notification.params?.childSessionId === 'provider-reuse-child',
  697. )).toEqual([{
  698. method: 'subagent.finished',
  699. params: {
  700. provider: 'reused-provider',
  701. agentId: 'provider-reuse-child',
  702. parentSessionId: 'provider-reuse-parent',
  703. childSessionId: 'provider-reuse-child',
  704. status: 'ok',
  705. stopReason: 'completed',
  706. lastAssistantMessage: [{ type: 'text', text: 'local' }],
  707. },
  708. }])
  709. await localRun.dispose()
  710. await remoteRun.dispose()
  711. unregisterRemote()
  712. await parent.dispose()
  713. await server.shutdown()
  714. } finally {
  715. await ctx.fiber.dispose()
  716. await rm(storageDir, { recursive: true, force: true })
  717. }
  718. })
  719. it('uses the recorded local flag when start was missed and ignores remote runs', async () => {
  720. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-subagent-fallback-'))
  721. const ctx = await makeHarness(storageDir)
  722. let parentHandle: AgentHandle | undefined
  723. let handle: AgentHandle | undefined
  724. let failedHandle: AgentHandle | undefined
  725. try {
  726. parentHandle = await ctx.agents.create({
  727. sessionId: SessionId('fallback-parent'),
  728. meta: { cwd: storageDir },
  729. agentOptions: { provider: 'deepseek-official', model: 'deepseek-official' },
  730. })
  731. handle = await parentHandle.agent.ctx.agents.create({
  732. sessionId: SessionId('fallback-child-session'),
  733. meta: { cwd: storageDir, parentSession: SessionId('fallback-parent') },
  734. agentOptions: { provider: 'deepseek-official', model: 'deepseek-official' },
  735. parentAgent: parentHandle.agent,
  736. })
  737. const fallbackChild = handle.agent
  738. failedHandle = await parentHandle.agent.ctx.agents.create({
  739. sessionId: SessionId('failed-child-session'),
  740. meta: { cwd: storageDir },
  741. agentOptions: { provider: 'deepseek-official', model: 'deepseek-official' },
  742. parentAgent: parentHandle.agent,
  743. })
  744. const missedStartResult = Promise.withResolvers<SubagentResult>()
  745. const disposeMissedStartProvider = ctx.subagents.registerProvider({
  746. name: 'fork',
  747. capabilities: { agentOptions: false, outputSchema: false, depthLimit: false, toolFilter: false, persona: false },
  748. inheritsParentContext: true,
  749. start: () => Promise.resolve({
  750. id: SessionId('fallback-child-session'),
  751. localAgent: fallbackChild,
  752. result: missedStartResult.promise,
  753. dispose: () => Promise.resolve(),
  754. }),
  755. })
  756. // Start before the server subscribes. The terminal payload still carries
  757. // this run's exact local child without reconstructing it from ids.
  758. const missedStartRun = await ctx.subagents.start('fork', {
  759. parent: parentHandle.agent,
  760. prompt: [],
  761. signal: new AbortController().signal,
  762. })
  763. const transport = new FakeTransport()
  764. const server = new HarnessSdkJsonRpcServer(ctx, transport, { maxTokensAsSuccess: true })
  765. missedStartResult.resolve({ output: [], stopReason: 'max-tokens' })
  766. await missedStartRun.result
  767. await Promise.resolve()
  768. await missedStartRun.dispose()
  769. disposeMissedStartProvider()
  770. // The server also missed this agent's creation but sees the exact child
  771. // on the run lifecycle payload.
  772. await settleSubagent(ctx, parentHandle.agent, {
  773. provider: 'fork-live-fallback',
  774. id: SessionId('fallback-child-session'),
  775. localAgent: fallbackChild,
  776. stopReason: 'completed',
  777. lastAssistantMessage: [],
  778. })
  779. await settleSubagent(ctx, parentHandle.agent, {
  780. provider: 'fork',
  781. id: SessionId('failed-child-session'),
  782. localAgent: failedHandle.agent,
  783. stopReason: 'error',
  784. })
  785. await settleSubagent(ctx, parentHandle.agent, {
  786. provider: 'fork',
  787. id: SessionId('missing-child-agent'),
  788. localAgent: undefined,
  789. stopReason: 'error',
  790. })
  791. // A result without output omits lastAssistantMessage from the wire; it
  792. // never sends `[]`.
  793. expect(transport.notifications).toContainEqual({
  794. method: 'subagent.finished',
  795. params: {
  796. provider: 'fork',
  797. agentId: 'fallback-child-session',
  798. parentSessionId: 'fallback-parent',
  799. childSessionId: 'fallback-child-session',
  800. status: 'ok',
  801. stopReason: 'max-tokens',
  802. },
  803. })
  804. expect(transport.notifications).toContainEqual({
  805. method: 'subagent.finished',
  806. params: {
  807. provider: 'fork',
  808. agentId: 'failed-child-session',
  809. parentSessionId: 'fallback-parent',
  810. childSessionId: 'failed-child-session',
  811. status: 'error',
  812. stopReason: 'error',
  813. },
  814. })
  815. expect(transport.notifications.some(n =>
  816. n.method === 'subagent.finished'
  817. && n.params?.agentId === 'missing-child-agent',
  818. )).toBe(false)
  819. await server.shutdown()
  820. } finally {
  821. await handle?.dispose()
  822. await failedHandle?.dispose()
  823. await parentHandle?.dispose()
  824. await ctx.fiber.dispose()
  825. await rm(storageDir, { recursive: true, force: true })
  826. }
  827. })
  828. it('does not re-register an LLM adapter whose provider already has an owner', async () => {
  829. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-existing-llm-'))
  830. const ctx = await makeHarness(storageDir)
  831. vi.stubEnv('DEEPSEEK_API_KEY', 'test-key')
  832. await ctx.plugin(LlmDeepSeek)
  833. try {
  834. const server = new HarnessSdkJsonRpcServer(ctx, new FakeTransport())
  835. const inspect = server as unknown as { hasAdapterFor(provider: string): boolean }
  836. expect(inspect.hasAdapterFor('deepseek-official')).toBe(true)
  837. expect(inspect.hasAdapterFor('missing-provider')).toBe(false)
  838. await server.initialize({ cwd: storageDir, provider: 'deepseek-official', model: 'preinstalled-model' })
  839. expect(ctx.get('llm')?.listProviders().filter(provider => provider.id === 'deepseek-official')).toEqual([{ id: 'deepseek-official', name: 'DeepSeek' }])
  840. await server.shutdown()
  841. } finally {
  842. await ctx.fiber.dispose()
  843. await rm(storageDir, { recursive: true, force: true })
  844. }
  845. })
  846. it('rejects a missing non-DeepSeek provider when an LLM service already exists', async () => {
  847. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-new-llm-'))
  848. const ctx = await makeHarness(storageDir)
  849. vi.stubEnv('DEEPSEEK_API_KEY', 'test-key')
  850. await ctx.plugin(LlmDeepSeek)
  851. try {
  852. const server = new HarnessSdkJsonRpcServer(ctx, new FakeTransport())
  853. await expect(server.initialize({ cwd: storageDir, provider: 'private', model: 'new-model' }))
  854. .rejects.toThrow('no adapter registered for provider "private"')
  855. expect(ctx.get('llm')?.listProviders()).toEqual([{ id: 'deepseek-official', name: 'DeepSeek' }])
  856. await server.shutdown()
  857. } finally {
  858. await ctx.fiber.dispose()
  859. await rm(storageDir, { recursive: true, force: true })
  860. }
  861. })
  862. it.each([0, -1, 1.5, Number.NaN, Number.MAX_SAFE_INTEGER + 1])(
  863. 'rejects invalid initialize maxTokens %s at the wire boundary',
  864. async (maxTokens) => {
  865. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-invalid-max-tokens-'))
  866. const ctx = await makeHarness(storageDir)
  867. try {
  868. const server = new HarnessSdkJsonRpcServer(ctx, new FakeTransport())
  869. await expect(server.initialize({
  870. cwd: storageDir,
  871. provider: 'deepseek-official',
  872. model: 'model',
  873. maxTokens,
  874. })).rejects.toThrow('initialize maxTokens must be a positive safe integer')
  875. await server.shutdown()
  876. } finally {
  877. await ctx.fiber.dispose()
  878. await rm(storageDir, { recursive: true, force: true })
  879. }
  880. },
  881. )
  882. it('rejects malformed initialize reasoningEffort values at the wire boundary', async () => {
  883. const ctx = new Context()
  884. const server = new HarnessSdkJsonRpcServer(ctx, new FakeTransport())
  885. try {
  886. for (const reasoningEffort of ['', 42]) {
  887. await expect(server.handleRequest('initialize', {
  888. cwd: '.',
  889. provider: 'deepseek-official',
  890. model: 'model',
  891. reasoningEffort,
  892. })).rejects.toThrow('initialize reasoningEffort must be a non-empty string')
  893. }
  894. await server.shutdown()
  895. } finally {
  896. await ctx.fiber.dispose()
  897. }
  898. })
  899. it('rejects an unavailable exact model during initialize', async () => {
  900. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-invalid-route-'))
  901. const ctx = await makeHarness(storageDir)
  902. class RejectingAdapter extends LlmAdapter {
  903. override resolveModel(provider: string, model: string): Promise<LlmResolvedModelInfo> {
  904. return Promise.reject(new Error(`model unavailable: ${provider}/${model}`))
  905. }
  906. async * stream(_options: GenerateOptions): AsyncIterable<StreamChunk> {
  907. throw new Error('unreachable')
  908. }
  909. }
  910. const disposeAdapter = ctx.llm.registerAdapter(['private'], new RejectingAdapter())
  911. try {
  912. const server = new HarnessSdkJsonRpcServer(ctx, new FakeTransport())
  913. await expect(server.initialize({ cwd: storageDir, provider: 'private', model: 'missing' }))
  914. .rejects.toThrow('model unavailable: private/missing')
  915. await expect(server.prompt({
  916. sessionId: 'invalid-route',
  917. contentBlocks: [{ type: 'text', text: 'must not run' }],
  918. })).rejects.toThrow('SDK server is not initialized')
  919. expect((server as unknown as { sessions: Map<string, unknown> }).sessions.size).toBe(0)
  920. await server.shutdown()
  921. } finally {
  922. disposeAdapter()
  923. await ctx.fiber.dispose()
  924. await rm(storageDir, { recursive: true, force: true })
  925. }
  926. })
  927. it('rejects prompts while exact-route initialization is pending', async () => {
  928. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-pending-route-'))
  929. const ctx = await makeHarness(storageDir)
  930. const resolution = Promise.withResolvers<LlmResolvedModelInfo>()
  931. const resolvedModel = { provider: 'private', id: 'selected', name: 'Selected' }
  932. let resolveModelCalled = false
  933. class PendingAdapter extends LlmAdapter {
  934. override resolveModel(): Promise<LlmResolvedModelInfo> {
  935. resolveModelCalled = true
  936. return resolution.promise
  937. }
  938. async * stream(_options: GenerateOptions): AsyncIterable<StreamChunk> {
  939. throw new Error('unreachable')
  940. }
  941. }
  942. const disposeAdapter = ctx.llm.registerAdapter(['private'], new PendingAdapter())
  943. try {
  944. const server = new HarnessSdkJsonRpcServer(ctx, new FakeTransport())
  945. const initialization = server.initialize({ cwd: storageDir, provider: 'private', model: 'selected' })
  946. await vi.waitFor(() => { expect(resolveModelCalled).toBe(true) })
  947. await expect(server.prompt({
  948. sessionId: 'too-early',
  949. contentBlocks: [{ type: 'text', text: 'must not run' }],
  950. })).rejects.toThrow('SDK server is not initialized')
  951. expect((server as unknown as { sessions: Map<string, unknown> }).sessions.size).toBe(0)
  952. resolution.resolve(resolvedModel)
  953. await initialization
  954. await server.shutdown()
  955. } finally {
  956. resolution.resolve(resolvedModel)
  957. disposeAdapter()
  958. await ctx.fiber.dispose()
  959. await rm(storageDir, { recursive: true, force: true })
  960. }
  961. })
  962. it('rejects an unsupported reasoning effort during initialize', async () => {
  963. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-unsupported-reasoning-'))
  964. const ctx = await makeHarness(storageDir)
  965. vi.stubEnv('DEEPSEEK_API_KEY', 'test-key')
  966. try {
  967. const server = new HarnessSdkJsonRpcServer(ctx, new FakeTransport())
  968. await expect(server.handleRequest('initialize', {
  969. cwd: storageDir,
  970. provider: 'deepseek-official',
  971. model: 'deepseek-v4-flash',
  972. reasoningEffort: 'impossible',
  973. })).rejects.toThrow('does not support reasoning effort "impossible"')
  974. expect((server as unknown as { sessions: Map<string, unknown> }).sessions.size).toBe(0)
  975. await server.shutdown()
  976. } finally {
  977. await ctx.fiber.dispose()
  978. await rm(storageDir, { recursive: true, force: true })
  979. }
  980. })
  981. it('reports no adapter when the LLM service is absent', async () => {
  982. const ctx = new Context()
  983. try {
  984. const server = new HarnessSdkJsonRpcServer(ctx, new FakeTransport()) as unknown as {
  985. hasAdapterFor(model: string): boolean
  986. shutdown(): Promise<Record<string, never>>
  987. }
  988. expect(server.hasAdapterFor('missing-model')).toBe(false)
  989. await server.shutdown()
  990. } finally {
  991. await ctx.fiber.dispose()
  992. }
  993. })
  994. it('rejects unknown JSON-RPC runtime methods', async () => {
  995. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-unknown-'))
  996. const ctx = await makeHarness(storageDir)
  997. try {
  998. const server = new HarnessSdkJsonRpcServer(ctx, new FakeTransport())
  999. await expect(server.handleRequest('does/not/exist', {}))
  1000. .rejects
  1001. .toThrow('unknown DeepSeek Harness SDK runtime method: does/not/exist')
  1002. await server.shutdown()
  1003. } finally {
  1004. await ctx.fiber.dispose()
  1005. await rm(storageDir, { recursive: true, force: true })
  1006. }
  1007. })
  1008. it('coalesces concurrent session creation and retries a failed creation', async () => {
  1009. let resolveShared: ((handle: AgentHandle) => void) | undefined
  1010. const sharedCreation = new Promise<AgentHandle>((resolve) => { resolveShared = resolve })
  1011. const sharedHandle = { agent: {} as Agent, dispose: vi.fn(() => Promise.resolve()) }
  1012. const retryHandle = { agent: {} as Agent, dispose: vi.fn(() => Promise.resolve()) }
  1013. const create = vi.fn<(options: unknown) => Promise<AgentHandle>>()
  1014. .mockReturnValueOnce(sharedCreation)
  1015. .mockRejectedValueOnce(new Error('creation failed'))
  1016. .mockResolvedValueOnce(retryHandle)
  1017. const ctx = {
  1018. on: vi.fn(() => () => undefined),
  1019. agents: { create, get: () => undefined },
  1020. get: () => undefined,
  1021. } as unknown as Context
  1022. const server = new HarnessSdkJsonRpcServer(ctx, new FakeTransport()) as unknown as {
  1023. getOrCreateSession(sessionId: string): Promise<{ handle: AgentHandle }>
  1024. shutdown(): Promise<Record<string, never>>
  1025. }
  1026. const first = server.getOrCreateSession('shared')
  1027. const second = server.getOrCreateSession('shared')
  1028. expect(create).toHaveBeenCalledTimes(1)
  1029. resolveShared?.(sharedHandle)
  1030. const [firstRecord, secondRecord] = await Promise.all([first, second])
  1031. expect(firstRecord).toBe(secondRecord)
  1032. await expect(server.getOrCreateSession('retry')).rejects.toThrow('creation failed')
  1033. await expect(server.getOrCreateSession('retry')).resolves.toMatchObject({ handle: retryHandle })
  1034. expect(create).toHaveBeenCalledTimes(3)
  1035. await server.shutdown()
  1036. expect(sharedHandle.dispose).toHaveBeenCalledOnce()
  1037. expect(retryHandle.dispose).toHaveBeenCalledOnce()
  1038. await expect(server.getOrCreateSession('after-shutdown')).rejects.toThrow('SDK server is shutting down')
  1039. })
  1040. it('resolves a relative cwd before creating the session', async () => {
  1041. const create = vi.fn<(options: unknown) => Promise<AgentHandle>>()
  1042. .mockResolvedValue({ agent: {} as Agent, dispose: () => Promise.resolve() })
  1043. const resolveCallConfig = vi.fn(async (config: unknown) => config)
  1044. const ctx = {
  1045. on: vi.fn(() => () => undefined),
  1046. agents: { create, get: () => undefined },
  1047. get: () => ({ listProviders: () => [{ id: 'mock', name: 'Mock' }], resolveCallConfig }),
  1048. } as unknown as Context
  1049. const server = new HarnessSdkJsonRpcServer(ctx, new FakeTransport()) as unknown as {
  1050. initialize(params: { cwd: string; provider: string; model: string; reasoningEffort?: string; maxTokens?: number }): Promise<unknown>
  1051. getOrCreateSession(sessionId: string): Promise<unknown>
  1052. shutdown(): Promise<Record<string, never>>
  1053. }
  1054. await server.initialize({ cwd: '.', provider: 'mock', model: 'model', reasoningEffort: 'high', maxTokens: 123 })
  1055. await server.getOrCreateSession('relative')
  1056. expect(resolveCallConfig).toHaveBeenCalledWith({
  1057. provider: 'mock',
  1058. model: 'model',
  1059. reasoningEffort: ReasoningEffortId('high'),
  1060. maxTokens: 123,
  1061. })
  1062. expect(create).toHaveBeenCalledWith(expect.objectContaining({
  1063. meta: { cwd: process.cwd() },
  1064. agentOptions: {
  1065. provider: 'mock',
  1066. model: 'model',
  1067. reasoningEffort: ReasoningEffortId('high'),
  1068. maxTokens: 123,
  1069. },
  1070. }))
  1071. await server.shutdown()
  1072. })
  1073. it('settles every teardown and aggregates multiple failures', async () => {
  1074. const firstDispose = vi.fn(() => { throw new Error('first teardown failed') })
  1075. const secondDispose = vi.fn(() => Promise.reject(new Error('second teardown failed')))
  1076. const ctx = {
  1077. on: vi.fn(() => () => undefined),
  1078. agents: { create: vi.fn(), get: () => undefined },
  1079. get: () => undefined,
  1080. } as unknown as Context
  1081. const server = new HarnessSdkJsonRpcServer(ctx, new FakeTransport()) as unknown as {
  1082. sessions: Map<string, { handle: AgentHandle; lastTurnEnd: undefined; activePrompt: boolean }>
  1083. shutdown(): Promise<Record<string, never>>
  1084. }
  1085. server.sessions.set('first', { handle: { agent: {} as Agent, dispose: firstDispose }, lastTurnEnd: undefined, activePrompt: false })
  1086. server.sessions.set('second', { handle: { agent: {} as Agent, dispose: secondDispose }, lastTurnEnd: undefined, activePrompt: false })
  1087. await expect(server.shutdown()).rejects.toThrow('SDK server teardown failed')
  1088. expect(firstDispose).toHaveBeenCalledOnce()
  1089. expect(secondDispose).toHaveBeenCalledOnce()
  1090. })
  1091. it('continues teardown after a subscription disposer fails', async () => {
  1092. let subscription = 0
  1093. const listenerFailure = new Error('listener teardown failed')
  1094. const on = vi.fn(() => {
  1095. subscription += 1
  1096. return subscription === 1 ? () => { throw listenerFailure } : () => undefined
  1097. })
  1098. const ctx = {
  1099. on,
  1100. agents: { create: vi.fn(), get: () => undefined },
  1101. get: () => undefined,
  1102. } as unknown as Context
  1103. const server = new HarnessSdkJsonRpcServer(ctx, new FakeTransport())
  1104. await expect(server.shutdown()).rejects.toBe(listenerFailure)
  1105. expect(on).toHaveBeenCalledTimes(4)
  1106. })
  1107. })