server.spec.ts 49 KB

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