server.spec.ts 42 KB

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