1
0

server.spec.ts 49 KB

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