server.spec.ts 44 KB

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