server.spec.ts 41 KB

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