server.spec.ts 39 KB

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