server.spec.ts 37 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929
  1. import { createServer } from 'node:http'
  2. import type { IncomingMessage, Server, ServerResponse } from 'node:http'
  3. import { mkdtemp, rm } from 'node:fs/promises'
  4. import { join } from 'node:path'
  5. import { tmpdir } from 'node:os'
  6. import { afterEach, describe, expect, it, vi } from 'vitest'
  7. import { Context } from 'cordis'
  8. import { type Agent, type AgentHandle } from '@deepseek-ai/dsh-agent'
  9. import { SessionId } from '@deepseek-ai/dsh-session'
  10. import * as agentCore from '@deepseek-ai/dsh-agent-spine-demo'
  11. import SessionPersistenceJsonl from '@deepseek-ai/dsh-session-persistence-jsonl'
  12. import * as LlmDeepSeek from '@deepseek-ai/dsh-llm-deepseek'
  13. import SubagentService, { type SubagentResult, type SubagentRunEndInfo } from '@deepseek-ai/dsh-subagent'
  14. import { HarnessSdkServer, type JsonRpcTransportPeer } from '../src/index.ts'
  15. class FakeTransport implements JsonRpcTransportPeer {
  16. notifications: { method: string; params?: Record<string, unknown> }[] = []
  17. async request(method: string, params: Record<string, unknown>): Promise<unknown> {
  18. throw new Error(`the SDK server should not call host JSON-RPC method ${method} with ${JSON.stringify(params)}`)
  19. }
  20. notify(method: string, params?: Record<string, unknown>): void {
  21. this.notifications.push(params === undefined ? { method } : { method, params })
  22. }
  23. }
  24. const servers: Server[] = []
  25. afterEach(async () => {
  26. await Promise.all(servers.splice(0).map(server => new Promise(resolve => server.close(resolve))))
  27. vi.unstubAllEnvs()
  28. })
  29. async function mockCompletionServer(): Promise<{ url: string; requests: unknown[]; headers: IncomingMessage['headers'][] }> {
  30. const requests: unknown[] = []
  31. const headers: IncomingMessage['headers'][] = []
  32. const server = createServer((request: IncomingMessage, response: ServerResponse) => {
  33. let body = ''
  34. request.on('data', (chunk: Buffer) => { body += chunk.toString('utf8') })
  35. request.on('end', () => {
  36. requests.push(JSON.parse(body))
  37. headers.push(request.headers)
  38. response.writeHead(200, { 'content-type': 'text/event-stream' })
  39. response.write('data: {"choices":[{"delta":{"role":"assistant","content":null,"reasoning_content":""}}]}\n\n')
  40. response.write('data: {"choices":[{"delta":{"content":"done"}}]}\n\n')
  41. response.write('data: {"choices":[{"delta":{"content":""},"finish_reason":"stop"}],"usage":{"prompt_tokens":3,"completion_tokens":1}}\n\n')
  42. response.write('data: [DONE]\n\n')
  43. response.end()
  44. })
  45. })
  46. servers.push(server)
  47. await new Promise<void>(resolve => server.listen(0, '127.0.0.1', resolve))
  48. const address = server.address()
  49. if (address === null || typeof address === 'string') throw new Error('no port')
  50. return { url: `http://127.0.0.1:${address.port}`, requests, headers }
  51. }
  52. async function makeHarness(storageDir: string) {
  53. const ctx = new Context()
  54. await ctx.plugin(agentCore, { workspaceContext: false })
  55. await ctx.plugin(SubagentService)
  56. await ctx.plugin(SessionPersistenceJsonl, { root: storageDir })
  57. await new Promise(resolve => setTimeout(resolve, 50))
  58. return ctx
  59. }
  60. /** Drive the owning service so test lifecycle events carry the real parent scope. */
  61. async function settleSubagent(
  62. ctx: Context,
  63. parent: Agent,
  64. info: Omit<SubagentRunEndInfo, 'runId' | 'local'> & { localAgent: Agent | undefined },
  65. beforeSettle?: () => Promise<void>,
  66. ): Promise<void> {
  67. const result = Promise.withResolvers<SubagentResult>()
  68. const disposeProvider = ctx.subagents.registerProvider({
  69. name: info.provider,
  70. capabilities: { outputSchema: false, depthLimit: false, toolFilter: false, persona: false },
  71. inheritsParentContext: false,
  72. async start() {
  73. return {
  74. id: info.id,
  75. localAgent: info.localAgent,
  76. result: result.promise,
  77. dispose: () => Promise.resolve(),
  78. }
  79. },
  80. })
  81. try {
  82. const run = await ctx.subagents.start(info.provider, {
  83. parent,
  84. prompt: [],
  85. signal: new AbortController().signal,
  86. })
  87. await beforeSettle?.()
  88. if (info.lastAssistantMessage === undefined) {
  89. result.reject(new Error('synthetic infrastructure failure'))
  90. } else {
  91. result.resolve({ output: info.lastAssistantMessage, stopReason: info.stopReason })
  92. }
  93. await run.result.then(() => undefined, () => undefined)
  94. await run.dispose()
  95. } finally {
  96. disposeProvider()
  97. }
  98. }
  99. describe('HarnessSdkServer', () => {
  100. it('creates a harness agent and calls the configured OpenAI-compatible endpoint', async () => {
  101. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-'))
  102. const llmServer = await mockCompletionServer()
  103. vi.stubEnv('DEEPSEEK_API_KEY', 'test-key')
  104. vi.stubEnv('DEEPSEEK_BASE_URL', llmServer.url)
  105. const ctx = await makeHarness(storageDir)
  106. try {
  107. const transport = new FakeTransport()
  108. const server = new HarnessSdkServer(ctx, transport)
  109. const init = await server.handleRequest('initialize', {
  110. cwd: storageDir,
  111. provider: 'deepseek',
  112. model: 'dsagent-model',
  113. }) as { serverInfo: { name: string } }
  114. expect(init.serverInfo.name).toBe('deepseek-harness-sdk-runtime')
  115. await server.handleRequest('session/prompt', {
  116. sessionId: 'main',
  117. contentBlocks: [{ type: 'text', text: 'fix it' }],
  118. })
  119. expect(llmServer.requests).toHaveLength(1)
  120. const body = llmServer.requests[0] as { model: string; messages: { role: string }[] }
  121. expect(body.model).toBe('dsagent-model')
  122. expect(body.messages[0]?.role).toBe('system')
  123. expect(body.messages.at(-1)?.role).toBe('user')
  124. expect(llmServer.headers[0]?.authorization).toBe('Bearer test-key')
  125. expect(transport.notifications.some(n => n.method === 'session.event')).toBe(true)
  126. expect(transport.notifications.at(-1)).toMatchObject({
  127. method: 'session.finished',
  128. params: { sessionId: 'main', status: 'ok' },
  129. })
  130. await server.handleRequest('session/prompt', {
  131. sessionId: 'main',
  132. contentBlocks: [{ type: 'text', text: 'again' }],
  133. })
  134. expect(llmServer.requests).toHaveLength(2)
  135. const orphanHandle = await ctx.agents.create({
  136. sessionId: SessionId('orphan-session'),
  137. meta: { cwd: storageDir },
  138. agentOptions: { provider: 'deepseek', model: 'dsagent-model' },
  139. })
  140. orphanHandle.agent.send([{ type: 'text', text: 'outside the sdk session map' }])
  141. await orphanHandle.agent.whenIdle()
  142. await orphanHandle.dispose()
  143. expect(llmServer.requests).toHaveLength(3)
  144. await server.handleRequest('shutdown', undefined)
  145. } finally {
  146. await ctx.fiber.dispose()
  147. await rm(storageDir, { recursive: true, force: true })
  148. }
  149. })
  150. it('rejects overlapping prompts for one session without serializing other sessions', async () => {
  151. let releaseMain: (() => void) | undefined
  152. const firstMainIdle = new Promise<void>((resolve) => { releaseMain = resolve })
  153. const mainWhenIdle = vi.fn<() => Promise<void>>()
  154. .mockReturnValueOnce(firstMainIdle)
  155. .mockResolvedValue(undefined)
  156. const mainSend = vi.fn()
  157. const mainAgent = {
  158. send: mainSend,
  159. whenIdle: mainWhenIdle,
  160. } as unknown as Agent
  161. const otherSend = vi.fn()
  162. const otherAgent = {
  163. send: otherSend,
  164. whenIdle: vi.fn(() => Promise.resolve()),
  165. } as unknown as Agent
  166. const mainHandle = { agent: mainAgent, dispose: vi.fn(() => Promise.resolve()) }
  167. const otherHandle = { agent: otherAgent, dispose: vi.fn(() => Promise.resolve()) }
  168. const create = vi.fn(async (options: { sessionId: SessionId }) =>
  169. String(options.sessionId) === 'main' ? mainHandle : otherHandle)
  170. const ctx = {
  171. on: vi.fn(() => () => undefined),
  172. agents: { create, get: () => undefined },
  173. get: () => undefined,
  174. } as unknown as Context
  175. const server = new HarnessSdkServer(ctx, new FakeTransport())
  176. const prompt = (sessionId: string, text: string) => server.prompt({
  177. sessionId,
  178. contentBlocks: [{ type: 'text', text }],
  179. })
  180. const first = prompt('main', 'first')
  181. await vi.waitFor(() => { expect(mainSend).toHaveBeenCalledOnce() })
  182. await expect(prompt('main', 'overlap')).rejects.toThrow('session already has an active prompt: main')
  183. await expect(prompt('other', 'independent')).resolves.toEqual({ accepted: true })
  184. releaseMain?.()
  185. await expect(first).resolves.toEqual({ accepted: true })
  186. await expect(prompt('main', 'sequential')).resolves.toEqual({ accepted: true })
  187. mainWhenIdle.mockRejectedValueOnce(new Error('turn wait failed'))
  188. await expect(prompt('main', 'failing')).rejects.toThrow('turn wait failed')
  189. await expect(prompt('main', 'after failure')).resolves.toEqual({ accepted: true })
  190. expect(mainSend).toHaveBeenCalledTimes(4)
  191. expect(otherSend).toHaveBeenCalledOnce()
  192. await server.shutdown()
  193. expect(mainHandle.dispose).toHaveBeenCalledOnce()
  194. expect(otherHandle.dispose).toHaveBeenCalledOnce()
  195. })
  196. it('notifies the host when a child session is created with parent lineage', async () => {
  197. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-subagent-'))
  198. const ctx = await makeHarness(storageDir)
  199. try {
  200. const transport = new FakeTransport()
  201. const server = new HarnessSdkServer(ctx, transport)
  202. ctx.sessions.create(SessionId('root-session'), {
  203. meta: { cwd: storageDir },
  204. })
  205. ctx.sessions.create(SessionId('child-session'), {
  206. meta: { cwd: storageDir, parentSession: SessionId('main') },
  207. })
  208. expect(transport.notifications).toContainEqual({
  209. method: 'subagent.started',
  210. params: {
  211. parentSessionId: 'main',
  212. childSessionId: 'child-session',
  213. },
  214. })
  215. await server.shutdown()
  216. } finally {
  217. await ctx.fiber.dispose()
  218. await rm(storageDir, { recursive: true, force: true })
  219. }
  220. })
  221. it('creates an SDK session without an optional system prompt', async () => {
  222. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-no-system-'))
  223. const llmServer = await mockCompletionServer()
  224. vi.stubEnv('DEEPSEEK_API_KEY', 'test-key')
  225. vi.stubEnv('DEEPSEEK_BASE_URL', llmServer.url)
  226. const ctx = await makeHarness(storageDir)
  227. try {
  228. const server = new HarnessSdkServer(ctx, new FakeTransport())
  229. await server.initialize({ cwd: storageDir, provider: 'deepseek', model: 'plain-model' })
  230. await server.prompt({
  231. sessionId: 'plain',
  232. contentBlocks: [{ type: 'text', text: 'hello' }],
  233. })
  234. expect(llmServer.requests).toHaveLength(1)
  235. await server.shutdown()
  236. } finally {
  237. await ctx.fiber.dispose()
  238. await rm(storageDir, { recursive: true, force: true })
  239. }
  240. })
  241. it('notifies the host when a subagent run settles', async () => {
  242. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-subagent-end-'))
  243. const ctx = await makeHarness(storageDir)
  244. try {
  245. const transport = new FakeTransport()
  246. const server = new HarnessSdkServer(ctx, transport)
  247. const parentHandle = await ctx.agents.create({
  248. sessionId: SessionId('main'),
  249. meta: { cwd: storageDir },
  250. agentOptions: { provider: 'deepseek', model: 'deepseek' },
  251. })
  252. // A custom in-process provider may own its child at the provider/root
  253. // scope while preserving durable parent lineage.
  254. const handle = await ctx.agents.create({
  255. sessionId: SessionId('child-session'),
  256. meta: { cwd: storageDir, parentSession: SessionId('main') },
  257. agentOptions: { provider: 'deepseek', model: 'deepseek' },
  258. })
  259. expect(ctx.agents.roots()).toContain(handle.agent)
  260. const parentlessHandle = await parentHandle.agent.ctx.agents.create({
  261. sessionId: SessionId('parentless-child-session'),
  262. meta: { cwd: storageDir },
  263. agentOptions: { model: 'deepseek' },
  264. })
  265. await settleSubagent(ctx, parentHandle.agent, {
  266. provider: 'spawn',
  267. id: SessionId('child-session'),
  268. localAgent: handle.agent,
  269. stopReason: 'completed',
  270. lastAssistantMessage: [{ type: 'text', text: 'child done' }],
  271. }, () => handle.dispose())
  272. await settleSubagent(ctx, parentHandle.agent, {
  273. provider: 'spawn',
  274. id: SessionId('parentless-child-session'),
  275. localAgent: parentlessHandle.agent,
  276. stopReason: 'error',
  277. }, () => parentlessHandle.dispose())
  278. expect(transport.notifications).toContainEqual({
  279. method: 'subagent.finished',
  280. params: {
  281. provider: 'spawn',
  282. agentId: 'child-session',
  283. parentSessionId: 'main',
  284. childSessionId: 'child-session',
  285. status: 'ok',
  286. stopReason: 'completed',
  287. lastAssistantMessage: [{ type: 'text', text: 'child done' }],
  288. },
  289. })
  290. expect(transport.notifications).toContainEqual({
  291. method: 'subagent.finished',
  292. params: {
  293. provider: 'spawn',
  294. agentId: 'parentless-child-session',
  295. parentSessionId: 'main',
  296. childSessionId: 'parentless-child-session',
  297. status: 'error',
  298. stopReason: 'error',
  299. },
  300. })
  301. await parentHandle.dispose()
  302. await server.shutdown()
  303. } finally {
  304. await ctx.fiber.dispose()
  305. await rm(storageDir, { recursive: true, force: true })
  306. }
  307. })
  308. it('ignores a remote run id that collides with a local child of the same parent', async () => {
  309. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-subagent-remote-collision-'))
  310. const ctx = await makeHarness(storageDir)
  311. try {
  312. const transport = new FakeTransport()
  313. const server = new HarnessSdkServer(ctx, transport)
  314. const parentHandle = await ctx.agents.create({
  315. sessionId: SessionId('collision-parent'),
  316. meta: { cwd: storageDir },
  317. agentOptions: { model: 'deepseek' },
  318. })
  319. const collidingChild = await parentHandle.agent.ctx.agents.create({
  320. sessionId: SessionId('remote-run-id'),
  321. meta: { cwd: storageDir, parentSession: SessionId('collision-parent') },
  322. agentOptions: { model: 'deepseek' },
  323. })
  324. await settleSubagent(ctx, parentHandle.agent, {
  325. provider: 'remote',
  326. id: SessionId('remote-run-id'),
  327. localAgent: undefined,
  328. stopReason: 'completed',
  329. lastAssistantMessage: [],
  330. })
  331. expect(transport.notifications.some(notification =>
  332. notification.method === 'subagent.finished'
  333. && notification.params?.agentId === 'remote-run-id',
  334. )).toBe(false)
  335. await collidingChild.dispose()
  336. await parentHandle.dispose()
  337. await server.shutdown()
  338. } finally {
  339. await ctx.fiber.dispose()
  340. await rm(storageDir, { recursive: true, force: true })
  341. }
  342. })
  343. it('retains locality across continuation runs on one live child', async () => {
  344. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-subagent-continuation-'))
  345. const ctx = await makeHarness(storageDir)
  346. try {
  347. const transport = new FakeTransport()
  348. const server = new HarnessSdkServer(ctx, transport)
  349. const parentHandle = await ctx.agents.create({
  350. sessionId: SessionId('continuation-parent'),
  351. meta: { cwd: storageDir },
  352. agentOptions: { model: 'deepseek' },
  353. })
  354. const childHandle = await parentHandle.agent.ctx.agents.create({
  355. sessionId: SessionId('continuation-child'),
  356. meta: { cwd: storageDir, parentSession: SessionId('continuation-parent') },
  357. agentOptions: { model: 'deepseek' },
  358. })
  359. await settleSubagent(ctx, parentHandle.agent, {
  360. provider: 'continuation',
  361. id: SessionId('continuation-child'),
  362. localAgent: childHandle.agent,
  363. stopReason: 'completed',
  364. lastAssistantMessage: [{ type: 'text', text: 'first' }],
  365. })
  366. await settleSubagent(ctx, parentHandle.agent, {
  367. provider: 'continuation',
  368. id: SessionId('continuation-child'),
  369. localAgent: childHandle.agent,
  370. stopReason: 'completed',
  371. lastAssistantMessage: [{ type: 'text', text: 'second' }],
  372. }, () => childHandle.dispose())
  373. expect(transport.notifications.filter(notification =>
  374. notification.method === 'subagent.finished'
  375. && notification.params?.childSessionId === 'continuation-child',
  376. )).toHaveLength(2)
  377. await parentHandle.dispose()
  378. await server.shutdown()
  379. } finally {
  380. await ctx.fiber.dispose()
  381. await rm(storageDir, { recursive: true, force: true })
  382. }
  383. })
  384. it('correlates reused local ids by parent scope when runs settle out of order', async () => {
  385. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-subagent-reuse-'))
  386. const ctx = await makeHarness(storageDir)
  387. try {
  388. const transport = new FakeTransport()
  389. const server = new HarnessSdkServer(ctx, transport)
  390. const oldParent = await ctx.agents.create({
  391. sessionId: SessionId('old-parent'),
  392. meta: { cwd: storageDir },
  393. agentOptions: { model: 'deepseek' },
  394. })
  395. const oldChild = await oldParent.agent.ctx.agents.create({
  396. sessionId: SessionId('reused-child'),
  397. meta: { cwd: storageDir, parentSession: SessionId('old-parent') },
  398. agentOptions: { model: 'deepseek' },
  399. })
  400. const first = Promise.withResolvers<SubagentResult>()
  401. const sameLifetime = Promise.withResolvers<SubagentResult>()
  402. const replacement = Promise.withResolvers<SubagentResult>()
  403. const results = [first.promise, sameLifetime.promise, replacement.promise]
  404. let starts = 0
  405. let currentLocalAgent = oldChild.agent
  406. const disposeProvider = ctx.subagents.registerProvider({
  407. name: 'reused',
  408. capabilities: { outputSchema: false, depthLimit: false, toolFilter: false, persona: false },
  409. inheritsParentContext: false,
  410. start() {
  411. const result = results[starts]
  412. starts += 1
  413. if (result === undefined) throw new Error('unexpected fourth reused-id run')
  414. return Promise.resolve({ id: SessionId('reused-child'), localAgent: currentLocalAgent, result, dispose: () => Promise.resolve() })
  415. },
  416. })
  417. const firstRun = await ctx.subagents.start('reused', {
  418. parent: oldParent.agent,
  419. prompt: [],
  420. signal: new AbortController().signal,
  421. })
  422. const sameLifetimeRun = await ctx.subagents.start('reused', {
  423. parent: oldParent.agent,
  424. prompt: [],
  425. signal: new AbortController().signal,
  426. })
  427. sameLifetime.resolve({ output: [{ type: 'text', text: 'same lifetime' }], stopReason: 'completed' })
  428. await sameLifetimeRun.result
  429. await oldChild.dispose()
  430. const newParent = await ctx.agents.create({
  431. sessionId: SessionId('new-parent'),
  432. meta: { cwd: storageDir },
  433. agentOptions: { model: 'deepseek' },
  434. })
  435. const newChild = await newParent.agent.ctx.agents.create({
  436. sessionId: SessionId('reused-child'),
  437. meta: { cwd: storageDir, parentSession: SessionId('new-parent') },
  438. agentOptions: { model: 'deepseek' },
  439. })
  440. currentLocalAgent = newChild.agent
  441. const secondRun = await ctx.subagents.start('reused', {
  442. parent: newParent.agent,
  443. prompt: [],
  444. signal: new AbortController().signal,
  445. })
  446. replacement.resolve({ output: [{ type: 'text', text: 'new lifetime' }], stopReason: 'completed' })
  447. await secondRun.result
  448. first.resolve({ output: [{ type: 'text', text: 'old lifetime' }], stopReason: 'completed' })
  449. await firstRun.result
  450. await Promise.resolve()
  451. const finished = transport.notifications.filter(notification =>
  452. notification.method === 'subagent.finished'
  453. && notification.params?.childSessionId === 'reused-child',
  454. )
  455. expect(finished.map(notification => notification.params?.lastAssistantMessage)).toEqual([
  456. [{ type: 'text', text: 'same lifetime' }],
  457. [{ type: 'text', text: 'new lifetime' }],
  458. [{ type: 'text', text: 'old lifetime' }],
  459. ])
  460. expect(finished.map(notification => notification.params?.parentSessionId)).toEqual([
  461. 'old-parent',
  462. 'new-parent',
  463. 'old-parent',
  464. ])
  465. await firstRun.dispose()
  466. await sameLifetimeRun.dispose()
  467. await secondRun.dispose()
  468. disposeProvider()
  469. await newChild.dispose()
  470. await oldParent.dispose()
  471. await newParent.dispose()
  472. await server.shutdown()
  473. } finally {
  474. await ctx.fiber.dispose()
  475. await rm(storageDir, { recursive: true, force: true })
  476. }
  477. })
  478. it('keeps locality bound to the accepted run across provider re-registration', async () => {
  479. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-subagent-provider-reuse-'))
  480. const ctx = await makeHarness(storageDir)
  481. try {
  482. const transport = new FakeTransport()
  483. const server = new HarnessSdkServer(ctx, transport)
  484. const parent = await ctx.agents.create({
  485. sessionId: SessionId('provider-reuse-parent'),
  486. meta: { cwd: storageDir },
  487. agentOptions: { model: 'deepseek' },
  488. })
  489. const child = await parent.agent.ctx.agents.create({
  490. sessionId: SessionId('provider-reuse-child'),
  491. meta: { cwd: storageDir, parentSession: SessionId('provider-reuse-parent') },
  492. agentOptions: { model: 'deepseek' },
  493. })
  494. const localResult = Promise.withResolvers<SubagentResult>()
  495. const remoteResult = Promise.withResolvers<SubagentResult>()
  496. const unregisterLocal = ctx.subagents.registerProvider({
  497. name: 'reused-provider',
  498. capabilities: { outputSchema: false, depthLimit: false, toolFilter: false, persona: false },
  499. inheritsParentContext: false,
  500. start: () => Promise.resolve({
  501. id: SessionId('provider-reuse-child'),
  502. localAgent: child.agent,
  503. result: localResult.promise,
  504. dispose: () => Promise.resolve(),
  505. }),
  506. })
  507. const localRun = await ctx.subagents.start('reused-provider', {
  508. parent: parent.agent,
  509. prompt: [],
  510. signal: new AbortController().signal,
  511. })
  512. unregisterLocal()
  513. const unregisterRemote = ctx.subagents.registerProvider({
  514. name: 'reused-provider',
  515. capabilities: { outputSchema: false, depthLimit: false, toolFilter: false, persona: false },
  516. inheritsParentContext: false,
  517. start: () => Promise.resolve({
  518. id: SessionId('provider-reuse-child'),
  519. localAgent: undefined,
  520. result: remoteResult.promise,
  521. dispose: () => Promise.resolve(),
  522. }),
  523. })
  524. const remoteRun = await ctx.subagents.start('reused-provider', {
  525. parent: parent.agent,
  526. prompt: [],
  527. signal: new AbortController().signal,
  528. })
  529. remoteResult.resolve({ output: [{ type: 'text', text: 'remote' }], stopReason: 'completed' })
  530. await remoteRun.result
  531. await Promise.resolve()
  532. expect(transport.notifications.some(notification =>
  533. notification.method === 'subagent.finished'
  534. && notification.params?.lastAssistantMessage !== undefined,
  535. )).toBe(false)
  536. await child.dispose()
  537. localResult.resolve({ output: [{ type: 'text', text: 'local' }], stopReason: 'completed' })
  538. await localRun.result
  539. await Promise.resolve()
  540. expect(transport.notifications.filter(notification =>
  541. notification.method === 'subagent.finished'
  542. && notification.params?.childSessionId === 'provider-reuse-child',
  543. )).toEqual([{
  544. method: 'subagent.finished',
  545. params: {
  546. provider: 'reused-provider',
  547. agentId: 'provider-reuse-child',
  548. parentSessionId: 'provider-reuse-parent',
  549. childSessionId: 'provider-reuse-child',
  550. status: 'ok',
  551. stopReason: 'completed',
  552. lastAssistantMessage: [{ type: 'text', text: 'local' }],
  553. },
  554. }])
  555. await localRun.dispose()
  556. await remoteRun.dispose()
  557. unregisterRemote()
  558. await parent.dispose()
  559. await server.shutdown()
  560. } finally {
  561. await ctx.fiber.dispose()
  562. await rm(storageDir, { recursive: true, force: true })
  563. }
  564. })
  565. it('uses explicit local provenance when start was missed and ignores remote runs', async () => {
  566. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-subagent-fallback-'))
  567. const ctx = await makeHarness(storageDir)
  568. let parentHandle: AgentHandle | undefined
  569. let handle: AgentHandle | undefined
  570. let failedHandle: AgentHandle | undefined
  571. try {
  572. parentHandle = await ctx.agents.create({
  573. sessionId: SessionId('fallback-parent'),
  574. meta: { cwd: storageDir },
  575. agentOptions: { provider: 'deepseek', model: 'deepseek' },
  576. })
  577. handle = await parentHandle.agent.ctx.agents.create({
  578. sessionId: SessionId('fallback-child-session'),
  579. meta: { cwd: storageDir, parentSession: SessionId('fallback-parent') },
  580. agentOptions: { provider: 'deepseek', model: 'deepseek' },
  581. })
  582. const fallbackChild = handle.agent
  583. failedHandle = await parentHandle.agent.ctx.agents.create({
  584. sessionId: SessionId('failed-child-session'),
  585. meta: { cwd: storageDir },
  586. agentOptions: { provider: 'deepseek', model: 'deepseek' },
  587. })
  588. const missedStartResult = Promise.withResolvers<SubagentResult>()
  589. const disposeMissedStartProvider = ctx.subagents.registerProvider({
  590. name: 'fork',
  591. capabilities: { outputSchema: false, depthLimit: false, toolFilter: false, persona: false },
  592. inheritsParentContext: true,
  593. start: () => Promise.resolve({
  594. id: SessionId('fallback-child-session'),
  595. localAgent: fallbackChild,
  596. result: missedStartResult.promise,
  597. dispose: () => Promise.resolve(),
  598. }),
  599. })
  600. // Start before the server subscribes. The terminal payload still carries
  601. // this run's exact local child without reconstructing it from ids.
  602. const missedStartRun = await ctx.subagents.start('fork', {
  603. parent: parentHandle.agent,
  604. prompt: [],
  605. signal: new AbortController().signal,
  606. })
  607. const transport = new FakeTransport()
  608. const server = new HarnessSdkServer(ctx, transport, { maxTokensAsSuccess: true })
  609. missedStartResult.resolve({ output: [], stopReason: 'max-tokens' })
  610. await missedStartRun.result
  611. await Promise.resolve()
  612. await missedStartRun.dispose()
  613. disposeMissedStartProvider()
  614. // The server also missed this agent's creation but sees the exact child
  615. // on the run lifecycle payload.
  616. await settleSubagent(ctx, parentHandle.agent, {
  617. provider: 'fork-live-fallback',
  618. id: SessionId('fallback-child-session'),
  619. localAgent: fallbackChild,
  620. stopReason: 'completed',
  621. lastAssistantMessage: [],
  622. })
  623. await settleSubagent(ctx, parentHandle.agent, {
  624. provider: 'fork',
  625. id: SessionId('failed-child-session'),
  626. localAgent: failedHandle.agent,
  627. stopReason: 'error',
  628. })
  629. await settleSubagent(ctx, parentHandle.agent, {
  630. provider: 'fork',
  631. id: SessionId('missing-child-agent'),
  632. localAgent: undefined,
  633. stopReason: 'error',
  634. })
  635. expect(transport.notifications).toContainEqual({
  636. method: 'subagent.finished',
  637. params: {
  638. provider: 'fork',
  639. agentId: 'fallback-child-session',
  640. parentSessionId: 'fallback-parent',
  641. childSessionId: 'fallback-child-session',
  642. status: 'ok',
  643. stopReason: 'max-tokens',
  644. lastAssistantMessage: [],
  645. },
  646. })
  647. expect(transport.notifications).toContainEqual({
  648. method: 'subagent.finished',
  649. params: {
  650. provider: 'fork',
  651. agentId: 'failed-child-session',
  652. parentSessionId: 'fallback-parent',
  653. childSessionId: 'failed-child-session',
  654. status: 'error',
  655. stopReason: 'error',
  656. },
  657. })
  658. expect(transport.notifications.some(n =>
  659. n.method === 'subagent.finished'
  660. && n.params?.agentId === 'missing-child-agent',
  661. )).toBe(false)
  662. await server.shutdown()
  663. } finally {
  664. await handle?.dispose()
  665. await failedHandle?.dispose()
  666. await parentHandle?.dispose()
  667. await ctx.fiber.dispose()
  668. await rm(storageDir, { recursive: true, force: true })
  669. }
  670. })
  671. it('does not re-register an LLM adapter whose provider already has an owner', async () => {
  672. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-existing-llm-'))
  673. const ctx = await makeHarness(storageDir)
  674. vi.stubEnv('DEEPSEEK_API_KEY', 'test-key')
  675. await ctx.plugin(LlmDeepSeek)
  676. try {
  677. const server = new HarnessSdkServer(ctx, new FakeTransport())
  678. const inspect = server as unknown as { hasAdapterFor(provider: string): boolean }
  679. expect(inspect.hasAdapterFor('deepseek')).toBe(true)
  680. expect(inspect.hasAdapterFor('missing-provider')).toBe(false)
  681. await server.initialize({ cwd: storageDir, provider: 'deepseek', model: 'preinstalled-model' })
  682. expect(ctx.get('llm')?.listProviders().filter(provider => provider.id === 'deepseek')).toEqual([{ id: 'deepseek', name: 'DeepSeek' }])
  683. await server.shutdown()
  684. } finally {
  685. await ctx.fiber.dispose()
  686. await rm(storageDir, { recursive: true, force: true })
  687. }
  688. })
  689. it('rejects a missing non-DeepSeek provider when an LLM service already exists', async () => {
  690. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-new-llm-'))
  691. const ctx = await makeHarness(storageDir)
  692. vi.stubEnv('DEEPSEEK_API_KEY', 'test-key')
  693. await ctx.plugin(LlmDeepSeek)
  694. try {
  695. const server = new HarnessSdkServer(ctx, new FakeTransport())
  696. await expect(server.initialize({ cwd: storageDir, provider: 'private', model: 'new-model' }))
  697. .rejects.toThrow('no adapter registered for provider "private"')
  698. expect(ctx.get('llm')?.listProviders()).toEqual([{ id: 'deepseek', name: 'DeepSeek' }])
  699. await server.shutdown()
  700. } finally {
  701. await ctx.fiber.dispose()
  702. await rm(storageDir, { recursive: true, force: true })
  703. }
  704. })
  705. it('classifies defensive finish states', async () => {
  706. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-finish-states-'))
  707. const ctx = await makeHarness(storageDir)
  708. try {
  709. const server = new HarnessSdkServer(ctx, new FakeTransport()) as unknown as {
  710. finishedStatus(reason: unknown): 'ok' | 'error'
  711. shutdown(): Promise<Record<string, never>>
  712. }
  713. expect(server.finishedStatus(undefined)).toBe('error')
  714. expect(server.finishedStatus({ kind: 'max-tokens' })).toBe('error')
  715. expect(server.finishedStatus({ kind: 'error' })).toBe('error')
  716. await server.shutdown()
  717. } finally {
  718. await ctx.fiber.dispose()
  719. await rm(storageDir, { recursive: true, force: true })
  720. }
  721. })
  722. it('can report max-token turn termination as an accepted evaluation result', async () => {
  723. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-max-tokens-success-'))
  724. const ctx = await makeHarness(storageDir)
  725. try {
  726. const server = new HarnessSdkServer(ctx, new FakeTransport(), { maxTokensAsSuccess: true }) as unknown as {
  727. finishedStatus(reason: unknown): 'ok' | 'error'
  728. shutdown(): Promise<Record<string, never>>
  729. }
  730. expect(server.finishedStatus({ kind: 'max-tokens' })).toBe('ok')
  731. expect(server.finishedStatus({ kind: 'error' })).toBe('error')
  732. await server.shutdown()
  733. } finally {
  734. await ctx.fiber.dispose()
  735. await rm(storageDir, { recursive: true, force: true })
  736. }
  737. })
  738. it('reports no adapter when the LLM service is absent', async () => {
  739. const ctx = new Context()
  740. try {
  741. const server = new HarnessSdkServer(ctx, new FakeTransport()) as unknown as {
  742. hasAdapterFor(model: string): boolean
  743. shutdown(): Promise<Record<string, never>>
  744. }
  745. expect(server.hasAdapterFor('missing-model')).toBe(false)
  746. await server.shutdown()
  747. } finally {
  748. await ctx.fiber.dispose()
  749. }
  750. })
  751. it('rejects unknown JSON-RPC runtime methods', async () => {
  752. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-unknown-'))
  753. const ctx = await makeHarness(storageDir)
  754. try {
  755. const server = new HarnessSdkServer(ctx, new FakeTransport())
  756. await expect(server.handleRequest('does/not/exist', {}))
  757. .rejects
  758. .toThrow('unknown DeepSeek Harness SDK runtime method: does/not/exist')
  759. await server.shutdown()
  760. } finally {
  761. await ctx.fiber.dispose()
  762. await rm(storageDir, { recursive: true, force: true })
  763. }
  764. })
  765. it('coalesces concurrent session creation and retries a failed creation', async () => {
  766. let resolveShared: ((handle: AgentHandle) => void) | undefined
  767. const sharedCreation = new Promise<AgentHandle>((resolve) => { resolveShared = resolve })
  768. const sharedHandle = { agent: {} as Agent, dispose: vi.fn(() => Promise.resolve()) }
  769. const retryHandle = { agent: {} as Agent, dispose: vi.fn(() => Promise.resolve()) }
  770. const create = vi.fn<(options: unknown) => Promise<AgentHandle>>()
  771. .mockReturnValueOnce(sharedCreation)
  772. .mockRejectedValueOnce(new Error('creation failed'))
  773. .mockResolvedValueOnce(retryHandle)
  774. const ctx = {
  775. on: vi.fn(() => () => undefined),
  776. agents: { create, get: () => undefined },
  777. get: () => undefined,
  778. } as unknown as Context
  779. const server = new HarnessSdkServer(ctx, new FakeTransport()) as unknown as {
  780. getOrCreateSession(sessionId: string): Promise<{ handle: AgentHandle }>
  781. shutdown(): Promise<Record<string, never>>
  782. }
  783. const first = server.getOrCreateSession('shared')
  784. const second = server.getOrCreateSession('shared')
  785. expect(create).toHaveBeenCalledTimes(1)
  786. resolveShared?.(sharedHandle)
  787. const [firstRecord, secondRecord] = await Promise.all([first, second])
  788. expect(firstRecord).toBe(secondRecord)
  789. await expect(server.getOrCreateSession('retry')).rejects.toThrow('creation failed')
  790. await expect(server.getOrCreateSession('retry')).resolves.toMatchObject({ handle: retryHandle })
  791. expect(create).toHaveBeenCalledTimes(3)
  792. await server.shutdown()
  793. expect(sharedHandle.dispose).toHaveBeenCalledOnce()
  794. expect(retryHandle.dispose).toHaveBeenCalledOnce()
  795. await expect(server.getOrCreateSession('after-shutdown')).rejects.toThrow('SDK server is shutting down')
  796. })
  797. it('resolves a relative cwd before creating the session', async () => {
  798. const create = vi.fn<(options: unknown) => Promise<AgentHandle>>()
  799. .mockResolvedValue({ agent: {} as Agent, dispose: () => Promise.resolve() })
  800. const ctx = {
  801. on: vi.fn(() => () => undefined),
  802. agents: { create, get: () => undefined },
  803. get: () => ({ listProviders: () => [{ id: 'mock', name: 'Mock' }] }),
  804. } as unknown as Context
  805. const server = new HarnessSdkServer(ctx, new FakeTransport()) as unknown as {
  806. initialize(params: { cwd: string; provider: string; model: string }): Promise<unknown>
  807. getOrCreateSession(sessionId: string): Promise<unknown>
  808. shutdown(): Promise<Record<string, never>>
  809. }
  810. await server.initialize({ cwd: '.', provider: 'mock', model: 'model' })
  811. await server.getOrCreateSession('relative')
  812. expect(create).toHaveBeenCalledWith(expect.objectContaining({ meta: { cwd: process.cwd() } }))
  813. await server.shutdown()
  814. })
  815. it('settles every teardown and aggregates multiple failures', async () => {
  816. const firstDispose = vi.fn(() => { throw new Error('first teardown failed') })
  817. const secondDispose = vi.fn(() => Promise.reject(new Error('second teardown failed')))
  818. const ctx = {
  819. on: vi.fn(() => () => undefined),
  820. agents: { create: vi.fn(), get: () => undefined },
  821. get: () => undefined,
  822. } as unknown as Context
  823. const server = new HarnessSdkServer(ctx, new FakeTransport()) as unknown as {
  824. sessions: Map<string, { handle: AgentHandle; lastTurnEnd: undefined; activePrompt: boolean }>
  825. shutdown(): Promise<Record<string, never>>
  826. }
  827. server.sessions.set('first', { handle: { agent: {} as Agent, dispose: firstDispose }, lastTurnEnd: undefined, activePrompt: false })
  828. server.sessions.set('second', { handle: { agent: {} as Agent, dispose: secondDispose }, lastTurnEnd: undefined, activePrompt: false })
  829. await expect(server.shutdown()).rejects.toThrow('SDK server teardown failed')
  830. expect(firstDispose).toHaveBeenCalledOnce()
  831. expect(secondDispose).toHaveBeenCalledOnce()
  832. })
  833. it('continues teardown after a subscription disposer fails', async () => {
  834. let subscription = 0
  835. const listenerFailure = new Error('listener teardown failed')
  836. const on = vi.fn(() => {
  837. subscription += 1
  838. return subscription === 1 ? () => { throw listenerFailure } : () => undefined
  839. })
  840. const ctx = {
  841. on,
  842. agents: { create: vi.fn(), get: () => undefined },
  843. get: () => undefined,
  844. } as unknown as Context
  845. const server = new HarnessSdkServer(ctx, new FakeTransport())
  846. await expect(server.shutdown()).rejects.toBe(listenerFailure)
  847. expect(on).toHaveBeenCalledTimes(3)
  848. })
  849. })