server.spec.ts 39 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985
  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({ content: [{ type: 'text', text: 'outside the sdk session map' }], source: { kind: 'user' } })
  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(input: { content: { type: 'text'; text: string }[]; source: { kind: 'user' } }) {
  209. session.append('turn/start', {
  210. turn: 1,
  211. trigger: { kind: 'message', source: input.source },
  212. })
  213. session.append('user/message', input, { surfaceOp: 'append' })
  214. session.append('turn/end', { turn: 1, reason: { kind: 'max-tokens' } })
  215. session.append('turn/start', {
  216. turn: 2,
  217. trigger: { kind: 'injection', source: { kind: 'plugin', plugin: 'late-metadata' } },
  218. })
  219. session.append('user/message', {
  220. content: [{ type: 'text', text: 'late metadata' }],
  221. source: { kind: 'plugin', plugin: 'late-metadata' },
  222. }, { surfaceOp: 'append' })
  223. session.append('turn/end', { turn: 2, reason: { kind: 'completed' } })
  224. return AgentMessageId('message-outcome')
  225. },
  226. whenIdle: () => Promise.resolve(),
  227. } satisfies Pick<Agent, 'session' | 'followup' | 'whenIdle'>) as unknown as Agent
  228. server.sessions.set('message-outcome', {
  229. handle: { agent, dispose: () => Promise.resolve() },
  230. lastTurnEnd: undefined,
  231. activePrompt: false,
  232. })
  233. await server.prompt({
  234. sessionId: 'message-outcome',
  235. contentBlocks: [{ type: 'text', text: 'bounded prompt' }],
  236. })
  237. expect(transport.notifications.findLast(notification => notification.method === 'session.finished'))
  238. .toEqual({
  239. method: 'session.finished',
  240. params: {
  241. sessionId: 'message-outcome',
  242. status: 'error',
  243. reason: { kind: 'max-tokens' },
  244. },
  245. })
  246. await server.shutdown()
  247. await ctx.fiber.dispose()
  248. })
  249. it('notifies the host when a child session is created with parent lineage', async () => {
  250. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-subagent-'))
  251. const ctx = await makeHarness(storageDir)
  252. try {
  253. const transport = new FakeTransport()
  254. const server = new HarnessSdkServer(ctx, transport)
  255. ctx.sessions.create(SessionId('root-session'), {
  256. meta: { cwd: storageDir },
  257. })
  258. ctx.sessions.create(SessionId('child-session'), {
  259. meta: { cwd: storageDir, parentSession: SessionId('main') },
  260. })
  261. expect(transport.notifications).toContainEqual({
  262. method: 'subagent.started',
  263. params: {
  264. parentSessionId: 'main',
  265. childSessionId: 'child-session',
  266. },
  267. })
  268. await server.shutdown()
  269. } finally {
  270. await ctx.fiber.dispose()
  271. await rm(storageDir, { recursive: true, force: true })
  272. }
  273. })
  274. it('creates an SDK session without an optional system prompt', async () => {
  275. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-no-system-'))
  276. const llmServer = await mockCompletionServer()
  277. vi.stubEnv('DEEPSEEK_API_KEY', 'test-key')
  278. vi.stubEnv('DEEPSEEK_BASE_URL', llmServer.url)
  279. const ctx = await makeHarness(storageDir)
  280. try {
  281. const server = new HarnessSdkServer(ctx, new FakeTransport())
  282. await server.initialize({ cwd: storageDir, provider: 'deepseek', model: 'plain-model' })
  283. await server.prompt({
  284. sessionId: 'plain',
  285. contentBlocks: [{ type: 'text', text: 'hello' }],
  286. })
  287. expect(llmServer.requests).toHaveLength(1)
  288. await server.shutdown()
  289. } finally {
  290. await ctx.fiber.dispose()
  291. await rm(storageDir, { recursive: true, force: true })
  292. }
  293. })
  294. it('notifies the host when a subagent run settles', async () => {
  295. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-subagent-end-'))
  296. const ctx = await makeHarness(storageDir)
  297. try {
  298. const transport = new FakeTransport()
  299. const server = new HarnessSdkServer(ctx, transport)
  300. const parentHandle = await ctx.agents.create({
  301. sessionId: SessionId('main'),
  302. meta: { cwd: storageDir },
  303. agentOptions: { provider: 'deepseek', model: 'deepseek' },
  304. })
  305. // A custom in-process provider may own its child at the provider/root
  306. // scope while preserving durable parent lineage.
  307. const handle = await ctx.agents.create({
  308. sessionId: SessionId('child-session'),
  309. meta: { cwd: storageDir, parentSession: SessionId('main') },
  310. agentOptions: { provider: 'deepseek', model: 'deepseek' },
  311. })
  312. expect(ctx.agents.roots()).toContain(handle.agent)
  313. const parentlessHandle = await parentHandle.agent.ctx.agents.create({
  314. sessionId: SessionId('parentless-child-session'),
  315. meta: { cwd: storageDir },
  316. agentOptions: { model: 'deepseek' },
  317. })
  318. await settleSubagent(ctx, parentHandle.agent, {
  319. provider: 'spawn',
  320. id: SessionId('child-session'),
  321. localAgent: handle.agent,
  322. stopReason: 'completed',
  323. lastAssistantMessage: [{ type: 'text', text: 'child done' }],
  324. }, () => handle.dispose())
  325. await settleSubagent(ctx, parentHandle.agent, {
  326. provider: 'spawn',
  327. id: SessionId('parentless-child-session'),
  328. localAgent: parentlessHandle.agent,
  329. stopReason: 'error',
  330. }, () => parentlessHandle.dispose())
  331. expect(transport.notifications).toContainEqual({
  332. method: 'subagent.finished',
  333. params: {
  334. provider: 'spawn',
  335. agentId: 'child-session',
  336. parentSessionId: 'main',
  337. childSessionId: 'child-session',
  338. status: 'ok',
  339. stopReason: 'completed',
  340. lastAssistantMessage: [{ type: 'text', text: 'child done' }],
  341. },
  342. })
  343. expect(transport.notifications).toContainEqual({
  344. method: 'subagent.finished',
  345. params: {
  346. provider: 'spawn',
  347. agentId: 'parentless-child-session',
  348. parentSessionId: 'main',
  349. childSessionId: 'parentless-child-session',
  350. status: 'error',
  351. stopReason: 'error',
  352. },
  353. })
  354. await parentHandle.dispose()
  355. await server.shutdown()
  356. } finally {
  357. await ctx.fiber.dispose()
  358. await rm(storageDir, { recursive: true, force: true })
  359. }
  360. })
  361. it('ignores a remote run id that collides with a local child of the same parent', async () => {
  362. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-subagent-remote-collision-'))
  363. const ctx = await makeHarness(storageDir)
  364. try {
  365. const transport = new FakeTransport()
  366. const server = new HarnessSdkServer(ctx, transport)
  367. const parentHandle = await ctx.agents.create({
  368. sessionId: SessionId('collision-parent'),
  369. meta: { cwd: storageDir },
  370. agentOptions: { model: 'deepseek' },
  371. })
  372. const collidingChild = await parentHandle.agent.ctx.agents.create({
  373. sessionId: SessionId('remote-run-id'),
  374. meta: { cwd: storageDir, parentSession: SessionId('collision-parent') },
  375. agentOptions: { model: 'deepseek' },
  376. })
  377. await settleSubagent(ctx, parentHandle.agent, {
  378. provider: 'remote',
  379. id: SessionId('remote-run-id'),
  380. localAgent: undefined,
  381. stopReason: 'completed',
  382. lastAssistantMessage: [],
  383. })
  384. expect(transport.notifications.some(notification =>
  385. notification.method === 'subagent.finished'
  386. && notification.params?.agentId === 'remote-run-id',
  387. )).toBe(false)
  388. await collidingChild.dispose()
  389. await parentHandle.dispose()
  390. await server.shutdown()
  391. } finally {
  392. await ctx.fiber.dispose()
  393. await rm(storageDir, { recursive: true, force: true })
  394. }
  395. })
  396. it('retains locality across continuation runs on one live child', async () => {
  397. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-subagent-continuation-'))
  398. const ctx = await makeHarness(storageDir)
  399. try {
  400. const transport = new FakeTransport()
  401. const server = new HarnessSdkServer(ctx, transport)
  402. const parentHandle = await ctx.agents.create({
  403. sessionId: SessionId('continuation-parent'),
  404. meta: { cwd: storageDir },
  405. agentOptions: { model: 'deepseek' },
  406. })
  407. const childHandle = await parentHandle.agent.ctx.agents.create({
  408. sessionId: SessionId('continuation-child'),
  409. meta: { cwd: storageDir, parentSession: SessionId('continuation-parent') },
  410. agentOptions: { model: 'deepseek' },
  411. })
  412. await settleSubagent(ctx, parentHandle.agent, {
  413. provider: 'continuation',
  414. id: SessionId('continuation-child'),
  415. localAgent: childHandle.agent,
  416. stopReason: 'completed',
  417. lastAssistantMessage: [{ type: 'text', text: 'first' }],
  418. })
  419. await settleSubagent(ctx, parentHandle.agent, {
  420. provider: 'continuation',
  421. id: SessionId('continuation-child'),
  422. localAgent: childHandle.agent,
  423. stopReason: 'completed',
  424. lastAssistantMessage: [{ type: 'text', text: 'second' }],
  425. }, () => childHandle.dispose())
  426. expect(transport.notifications.filter(notification =>
  427. notification.method === 'subagent.finished'
  428. && notification.params?.childSessionId === 'continuation-child',
  429. )).toHaveLength(2)
  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('correlates reused local ids by parent scope when runs settle out of order', async () => {
  438. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-subagent-reuse-'))
  439. const ctx = await makeHarness(storageDir)
  440. try {
  441. const transport = new FakeTransport()
  442. const server = new HarnessSdkServer(ctx, transport)
  443. const oldParent = await ctx.agents.create({
  444. sessionId: SessionId('old-parent'),
  445. meta: { cwd: storageDir },
  446. agentOptions: { model: 'deepseek' },
  447. })
  448. const oldChild = await oldParent.agent.ctx.agents.create({
  449. sessionId: SessionId('reused-child'),
  450. meta: { cwd: storageDir, parentSession: SessionId('old-parent') },
  451. agentOptions: { model: 'deepseek' },
  452. })
  453. const first = Promise.withResolvers<SubagentResult>()
  454. const sameLifetime = Promise.withResolvers<SubagentResult>()
  455. const replacement = Promise.withResolvers<SubagentResult>()
  456. const results = [first.promise, sameLifetime.promise, replacement.promise]
  457. let starts = 0
  458. let currentLocalAgent = oldChild.agent
  459. const disposeProvider = ctx.subagents.registerProvider({
  460. name: 'reused',
  461. capabilities: { outputSchema: false, depthLimit: false, toolFilter: false, persona: false },
  462. inheritsParentContext: false,
  463. start() {
  464. const result = results[starts]
  465. starts += 1
  466. if (result === undefined) throw new Error('unexpected fourth reused-id run')
  467. return Promise.resolve({ id: SessionId('reused-child'), localAgent: currentLocalAgent, result, dispose: () => Promise.resolve() })
  468. },
  469. })
  470. const firstRun = await ctx.subagents.start('reused', {
  471. parent: oldParent.agent,
  472. prompt: [],
  473. signal: new AbortController().signal,
  474. })
  475. const sameLifetimeRun = await ctx.subagents.start('reused', {
  476. parent: oldParent.agent,
  477. prompt: [],
  478. signal: new AbortController().signal,
  479. })
  480. sameLifetime.resolve({ output: [{ type: 'text', text: 'same lifetime' }], stopReason: 'completed' })
  481. await sameLifetimeRun.result
  482. await oldChild.dispose()
  483. const newParent = await ctx.agents.create({
  484. sessionId: SessionId('new-parent'),
  485. meta: { cwd: storageDir },
  486. agentOptions: { model: 'deepseek' },
  487. })
  488. const newChild = await newParent.agent.ctx.agents.create({
  489. sessionId: SessionId('reused-child'),
  490. meta: { cwd: storageDir, parentSession: SessionId('new-parent') },
  491. agentOptions: { model: 'deepseek' },
  492. })
  493. currentLocalAgent = newChild.agent
  494. const secondRun = await ctx.subagents.start('reused', {
  495. parent: newParent.agent,
  496. prompt: [],
  497. signal: new AbortController().signal,
  498. })
  499. replacement.resolve({ output: [{ type: 'text', text: 'new lifetime' }], stopReason: 'completed' })
  500. await secondRun.result
  501. first.resolve({ output: [{ type: 'text', text: 'old lifetime' }], stopReason: 'completed' })
  502. await firstRun.result
  503. await Promise.resolve()
  504. const finished = transport.notifications.filter(notification =>
  505. notification.method === 'subagent.finished'
  506. && notification.params?.childSessionId === 'reused-child',
  507. )
  508. expect(finished.map(notification => notification.params?.lastAssistantMessage)).toEqual([
  509. [{ type: 'text', text: 'same lifetime' }],
  510. [{ type: 'text', text: 'new lifetime' }],
  511. [{ type: 'text', text: 'old lifetime' }],
  512. ])
  513. expect(finished.map(notification => notification.params?.parentSessionId)).toEqual([
  514. 'old-parent',
  515. 'new-parent',
  516. 'old-parent',
  517. ])
  518. await firstRun.dispose()
  519. await sameLifetimeRun.dispose()
  520. await secondRun.dispose()
  521. disposeProvider()
  522. await newChild.dispose()
  523. await oldParent.dispose()
  524. await newParent.dispose()
  525. await server.shutdown()
  526. } finally {
  527. await ctx.fiber.dispose()
  528. await rm(storageDir, { recursive: true, force: true })
  529. }
  530. })
  531. it('keeps locality bound to the accepted run across provider re-registration', async () => {
  532. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-subagent-provider-reuse-'))
  533. const ctx = await makeHarness(storageDir)
  534. try {
  535. const transport = new FakeTransport()
  536. const server = new HarnessSdkServer(ctx, transport)
  537. const parent = await ctx.agents.create({
  538. sessionId: SessionId('provider-reuse-parent'),
  539. meta: { cwd: storageDir },
  540. agentOptions: { model: 'deepseek' },
  541. })
  542. const child = await parent.agent.ctx.agents.create({
  543. sessionId: SessionId('provider-reuse-child'),
  544. meta: { cwd: storageDir, parentSession: SessionId('provider-reuse-parent') },
  545. agentOptions: { model: 'deepseek' },
  546. })
  547. const localResult = Promise.withResolvers<SubagentResult>()
  548. const remoteResult = Promise.withResolvers<SubagentResult>()
  549. const unregisterLocal = ctx.subagents.registerProvider({
  550. name: 'reused-provider',
  551. capabilities: { outputSchema: false, depthLimit: false, toolFilter: false, persona: false },
  552. inheritsParentContext: false,
  553. start: () => Promise.resolve({
  554. id: SessionId('provider-reuse-child'),
  555. localAgent: child.agent,
  556. result: localResult.promise,
  557. dispose: () => Promise.resolve(),
  558. }),
  559. })
  560. const localRun = await ctx.subagents.start('reused-provider', {
  561. parent: parent.agent,
  562. prompt: [],
  563. signal: new AbortController().signal,
  564. })
  565. unregisterLocal()
  566. const unregisterRemote = ctx.subagents.registerProvider({
  567. name: 'reused-provider',
  568. capabilities: { outputSchema: false, depthLimit: false, toolFilter: false, persona: false },
  569. inheritsParentContext: false,
  570. start: () => Promise.resolve({
  571. id: SessionId('provider-reuse-child'),
  572. localAgent: undefined,
  573. result: remoteResult.promise,
  574. dispose: () => Promise.resolve(),
  575. }),
  576. })
  577. const remoteRun = await ctx.subagents.start('reused-provider', {
  578. parent: parent.agent,
  579. prompt: [],
  580. signal: new AbortController().signal,
  581. })
  582. remoteResult.resolve({ output: [{ type: 'text', text: 'remote' }], stopReason: 'completed' })
  583. await remoteRun.result
  584. await Promise.resolve()
  585. expect(transport.notifications.some(notification =>
  586. notification.method === 'subagent.finished'
  587. && notification.params?.lastAssistantMessage !== undefined,
  588. )).toBe(false)
  589. await child.dispose()
  590. localResult.resolve({ output: [{ type: 'text', text: 'local' }], stopReason: 'completed' })
  591. await localRun.result
  592. await Promise.resolve()
  593. expect(transport.notifications.filter(notification =>
  594. notification.method === 'subagent.finished'
  595. && notification.params?.childSessionId === 'provider-reuse-child',
  596. )).toEqual([{
  597. method: 'subagent.finished',
  598. params: {
  599. provider: 'reused-provider',
  600. agentId: 'provider-reuse-child',
  601. parentSessionId: 'provider-reuse-parent',
  602. childSessionId: 'provider-reuse-child',
  603. status: 'ok',
  604. stopReason: 'completed',
  605. lastAssistantMessage: [{ type: 'text', text: 'local' }],
  606. },
  607. }])
  608. await localRun.dispose()
  609. await remoteRun.dispose()
  610. unregisterRemote()
  611. await parent.dispose()
  612. await server.shutdown()
  613. } finally {
  614. await ctx.fiber.dispose()
  615. await rm(storageDir, { recursive: true, force: true })
  616. }
  617. })
  618. it('uses explicit local provenance when start was missed and ignores remote runs', async () => {
  619. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-subagent-fallback-'))
  620. const ctx = await makeHarness(storageDir)
  621. let parentHandle: AgentHandle | undefined
  622. let handle: AgentHandle | undefined
  623. let failedHandle: AgentHandle | undefined
  624. try {
  625. parentHandle = await ctx.agents.create({
  626. sessionId: SessionId('fallback-parent'),
  627. meta: { cwd: storageDir },
  628. agentOptions: { provider: 'deepseek', model: 'deepseek' },
  629. })
  630. handle = await parentHandle.agent.ctx.agents.create({
  631. sessionId: SessionId('fallback-child-session'),
  632. meta: { cwd: storageDir, parentSession: SessionId('fallback-parent') },
  633. agentOptions: { provider: 'deepseek', model: 'deepseek' },
  634. })
  635. const fallbackChild = handle.agent
  636. failedHandle = await parentHandle.agent.ctx.agents.create({
  637. sessionId: SessionId('failed-child-session'),
  638. meta: { cwd: storageDir },
  639. agentOptions: { provider: 'deepseek', model: 'deepseek' },
  640. })
  641. const missedStartResult = Promise.withResolvers<SubagentResult>()
  642. const disposeMissedStartProvider = ctx.subagents.registerProvider({
  643. name: 'fork',
  644. capabilities: { outputSchema: false, depthLimit: false, toolFilter: false, persona: false },
  645. inheritsParentContext: true,
  646. start: () => Promise.resolve({
  647. id: SessionId('fallback-child-session'),
  648. localAgent: fallbackChild,
  649. result: missedStartResult.promise,
  650. dispose: () => Promise.resolve(),
  651. }),
  652. })
  653. // Start before the server subscribes. The terminal payload still carries
  654. // this run's exact local child without reconstructing it from ids.
  655. const missedStartRun = await ctx.subagents.start('fork', {
  656. parent: parentHandle.agent,
  657. prompt: [],
  658. signal: new AbortController().signal,
  659. })
  660. const transport = new FakeTransport()
  661. const server = new HarnessSdkServer(ctx, transport, { maxTokensAsSuccess: true })
  662. missedStartResult.resolve({ output: [], stopReason: 'max-tokens' })
  663. await missedStartRun.result
  664. await Promise.resolve()
  665. await missedStartRun.dispose()
  666. disposeMissedStartProvider()
  667. // The server also missed this agent's creation but sees the exact child
  668. // on the run lifecycle payload.
  669. await settleSubagent(ctx, parentHandle.agent, {
  670. provider: 'fork-live-fallback',
  671. id: SessionId('fallback-child-session'),
  672. localAgent: fallbackChild,
  673. stopReason: 'completed',
  674. lastAssistantMessage: [],
  675. })
  676. await settleSubagent(ctx, parentHandle.agent, {
  677. provider: 'fork',
  678. id: SessionId('failed-child-session'),
  679. localAgent: failedHandle.agent,
  680. stopReason: 'error',
  681. })
  682. await settleSubagent(ctx, parentHandle.agent, {
  683. provider: 'fork',
  684. id: SessionId('missing-child-agent'),
  685. localAgent: undefined,
  686. stopReason: 'error',
  687. })
  688. expect(transport.notifications).toContainEqual({
  689. method: 'subagent.finished',
  690. params: {
  691. provider: 'fork',
  692. agentId: 'fallback-child-session',
  693. parentSessionId: 'fallback-parent',
  694. childSessionId: 'fallback-child-session',
  695. status: 'ok',
  696. stopReason: 'max-tokens',
  697. lastAssistantMessage: [],
  698. },
  699. })
  700. expect(transport.notifications).toContainEqual({
  701. method: 'subagent.finished',
  702. params: {
  703. provider: 'fork',
  704. agentId: 'failed-child-session',
  705. parentSessionId: 'fallback-parent',
  706. childSessionId: 'failed-child-session',
  707. status: 'error',
  708. stopReason: 'error',
  709. },
  710. })
  711. expect(transport.notifications.some(n =>
  712. n.method === 'subagent.finished'
  713. && n.params?.agentId === 'missing-child-agent',
  714. )).toBe(false)
  715. await server.shutdown()
  716. } finally {
  717. await handle?.dispose()
  718. await failedHandle?.dispose()
  719. await parentHandle?.dispose()
  720. await ctx.fiber.dispose()
  721. await rm(storageDir, { recursive: true, force: true })
  722. }
  723. })
  724. it('does not re-register an LLM adapter whose provider already has an owner', async () => {
  725. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-existing-llm-'))
  726. const ctx = await makeHarness(storageDir)
  727. vi.stubEnv('DEEPSEEK_API_KEY', 'test-key')
  728. await ctx.plugin(LlmDeepSeek)
  729. try {
  730. const server = new HarnessSdkServer(ctx, new FakeTransport())
  731. const inspect = server as unknown as { hasAdapterFor(provider: string): boolean }
  732. expect(inspect.hasAdapterFor('deepseek')).toBe(true)
  733. expect(inspect.hasAdapterFor('missing-provider')).toBe(false)
  734. await server.initialize({ cwd: storageDir, provider: 'deepseek', model: 'preinstalled-model' })
  735. expect(ctx.get('llm')?.listProviders().filter(provider => provider.id === 'deepseek')).toEqual([{ id: 'deepseek', name: 'DeepSeek' }])
  736. await server.shutdown()
  737. } finally {
  738. await ctx.fiber.dispose()
  739. await rm(storageDir, { recursive: true, force: true })
  740. }
  741. })
  742. it('rejects a missing non-DeepSeek provider when an LLM service already exists', async () => {
  743. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-new-llm-'))
  744. const ctx = await makeHarness(storageDir)
  745. vi.stubEnv('DEEPSEEK_API_KEY', 'test-key')
  746. await ctx.plugin(LlmDeepSeek)
  747. try {
  748. const server = new HarnessSdkServer(ctx, new FakeTransport())
  749. await expect(server.initialize({ cwd: storageDir, provider: 'private', model: 'new-model' }))
  750. .rejects.toThrow('no adapter registered for provider "private"')
  751. expect(ctx.get('llm')?.listProviders()).toEqual([{ id: 'deepseek', name: 'DeepSeek' }])
  752. await server.shutdown()
  753. } finally {
  754. await ctx.fiber.dispose()
  755. await rm(storageDir, { recursive: true, force: true })
  756. }
  757. })
  758. it('classifies defensive finish states', async () => {
  759. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-finish-states-'))
  760. const ctx = await makeHarness(storageDir)
  761. try {
  762. const server = new HarnessSdkServer(ctx, new FakeTransport()) as unknown as {
  763. finishedStatus(reason: unknown): 'ok' | 'error'
  764. shutdown(): Promise<Record<string, never>>
  765. }
  766. expect(server.finishedStatus(undefined)).toBe('error')
  767. expect(server.finishedStatus({ kind: 'max-tokens' })).toBe('error')
  768. expect(server.finishedStatus({ kind: 'error' })).toBe('error')
  769. await server.shutdown()
  770. } finally {
  771. await ctx.fiber.dispose()
  772. await rm(storageDir, { recursive: true, force: true })
  773. }
  774. })
  775. it('can report max-token turn termination as an accepted evaluation result', async () => {
  776. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-max-tokens-success-'))
  777. const ctx = await makeHarness(storageDir)
  778. try {
  779. const server = new HarnessSdkServer(ctx, new FakeTransport(), { maxTokensAsSuccess: true }) as unknown as {
  780. finishedStatus(reason: unknown): 'ok' | 'error'
  781. shutdown(): Promise<Record<string, never>>
  782. }
  783. expect(server.finishedStatus({ kind: 'max-tokens' })).toBe('ok')
  784. expect(server.finishedStatus({ kind: 'error' })).toBe('error')
  785. await server.shutdown()
  786. } finally {
  787. await ctx.fiber.dispose()
  788. await rm(storageDir, { recursive: true, force: true })
  789. }
  790. })
  791. it('reports no adapter when the LLM service is absent', async () => {
  792. const ctx = new Context()
  793. try {
  794. const server = new HarnessSdkServer(ctx, new FakeTransport()) as unknown as {
  795. hasAdapterFor(model: string): boolean
  796. shutdown(): Promise<Record<string, never>>
  797. }
  798. expect(server.hasAdapterFor('missing-model')).toBe(false)
  799. await server.shutdown()
  800. } finally {
  801. await ctx.fiber.dispose()
  802. }
  803. })
  804. it('rejects unknown JSON-RPC runtime methods', async () => {
  805. const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-unknown-'))
  806. const ctx = await makeHarness(storageDir)
  807. try {
  808. const server = new HarnessSdkServer(ctx, new FakeTransport())
  809. await expect(server.handleRequest('does/not/exist', {}))
  810. .rejects
  811. .toThrow('unknown DeepSeek Harness SDK runtime method: does/not/exist')
  812. await server.shutdown()
  813. } finally {
  814. await ctx.fiber.dispose()
  815. await rm(storageDir, { recursive: true, force: true })
  816. }
  817. })
  818. it('coalesces concurrent session creation and retries a failed creation', async () => {
  819. let resolveShared: ((handle: AgentHandle) => void) | undefined
  820. const sharedCreation = new Promise<AgentHandle>((resolve) => { resolveShared = resolve })
  821. const sharedHandle = { agent: {} as Agent, dispose: vi.fn(() => Promise.resolve()) }
  822. const retryHandle = { agent: {} as Agent, dispose: vi.fn(() => Promise.resolve()) }
  823. const create = vi.fn<(options: unknown) => Promise<AgentHandle>>()
  824. .mockReturnValueOnce(sharedCreation)
  825. .mockRejectedValueOnce(new Error('creation failed'))
  826. .mockResolvedValueOnce(retryHandle)
  827. const ctx = {
  828. on: vi.fn(() => () => undefined),
  829. agents: { create, get: () => undefined },
  830. get: () => undefined,
  831. } as unknown as Context
  832. const server = new HarnessSdkServer(ctx, new FakeTransport()) as unknown as {
  833. getOrCreateSession(sessionId: string): Promise<{ handle: AgentHandle }>
  834. shutdown(): Promise<Record<string, never>>
  835. }
  836. const first = server.getOrCreateSession('shared')
  837. const second = server.getOrCreateSession('shared')
  838. expect(create).toHaveBeenCalledTimes(1)
  839. resolveShared?.(sharedHandle)
  840. const [firstRecord, secondRecord] = await Promise.all([first, second])
  841. expect(firstRecord).toBe(secondRecord)
  842. await expect(server.getOrCreateSession('retry')).rejects.toThrow('creation failed')
  843. await expect(server.getOrCreateSession('retry')).resolves.toMatchObject({ handle: retryHandle })
  844. expect(create).toHaveBeenCalledTimes(3)
  845. await server.shutdown()
  846. expect(sharedHandle.dispose).toHaveBeenCalledOnce()
  847. expect(retryHandle.dispose).toHaveBeenCalledOnce()
  848. await expect(server.getOrCreateSession('after-shutdown')).rejects.toThrow('SDK server is shutting down')
  849. })
  850. it('resolves a relative cwd before creating the session', async () => {
  851. const create = vi.fn<(options: unknown) => Promise<AgentHandle>>()
  852. .mockResolvedValue({ agent: {} as Agent, dispose: () => Promise.resolve() })
  853. const ctx = {
  854. on: vi.fn(() => () => undefined),
  855. agents: { create, get: () => undefined },
  856. get: () => ({ listProviders: () => [{ id: 'mock', name: 'Mock' }] }),
  857. } as unknown as Context
  858. const server = new HarnessSdkServer(ctx, new FakeTransport()) as unknown as {
  859. initialize(params: { cwd: string; provider: string; model: string }): Promise<unknown>
  860. getOrCreateSession(sessionId: string): Promise<unknown>
  861. shutdown(): Promise<Record<string, never>>
  862. }
  863. await server.initialize({ cwd: '.', provider: 'mock', model: 'model' })
  864. await server.getOrCreateSession('relative')
  865. expect(create).toHaveBeenCalledWith(expect.objectContaining({ meta: { cwd: process.cwd() } }))
  866. await server.shutdown()
  867. })
  868. it('settles every teardown and aggregates multiple failures', async () => {
  869. const firstDispose = vi.fn(() => { throw new Error('first teardown failed') })
  870. const secondDispose = vi.fn(() => Promise.reject(new Error('second teardown failed')))
  871. const ctx = {
  872. on: vi.fn(() => () => undefined),
  873. agents: { create: vi.fn(), get: () => undefined },
  874. get: () => undefined,
  875. } as unknown as Context
  876. const server = new HarnessSdkServer(ctx, new FakeTransport()) as unknown as {
  877. sessions: Map<string, { handle: AgentHandle; lastTurnEnd: undefined; activePrompt: boolean }>
  878. shutdown(): Promise<Record<string, never>>
  879. }
  880. server.sessions.set('first', { handle: { agent: {} as Agent, dispose: firstDispose }, lastTurnEnd: undefined, activePrompt: false })
  881. server.sessions.set('second', { handle: { agent: {} as Agent, dispose: secondDispose }, lastTurnEnd: undefined, activePrompt: false })
  882. await expect(server.shutdown()).rejects.toThrow('SDK server teardown failed')
  883. expect(firstDispose).toHaveBeenCalledOnce()
  884. expect(secondDispose).toHaveBeenCalledOnce()
  885. })
  886. it('continues teardown after a subscription disposer fails', async () => {
  887. let subscription = 0
  888. const listenerFailure = new Error('listener teardown failed')
  889. const on = vi.fn(() => {
  890. subscription += 1
  891. return subscription === 1 ? () => { throw listenerFailure } : () => undefined
  892. })
  893. const ctx = {
  894. on,
  895. agents: { create: vi.fn(), get: () => undefined },
  896. get: () => undefined,
  897. } as unknown as Context
  898. const server = new HarnessSdkServer(ctx, new FakeTransport())
  899. await expect(server.shutdown()).rejects.toBe(listenerFailure)
  900. expect(on).toHaveBeenCalledTimes(3)
  901. })
  902. })