server.spec.ts 48 KB

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