server.spec.ts 39 KB

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