host-runtime.spec.ts 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367
  1. import { mkdtempSync } from 'node:fs'
  2. import { tmpdir } from 'node:os'
  3. import { join } from 'node:path'
  4. import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
  5. import type { Context } from 'cordis'
  6. import type { Agent } from '@deepseek-ai/dsh-agent'
  7. import { agentEvents } from '@deepseek-ai/dsh-agent'
  8. import type { GenerateOptions, StreamChunk } from '@deepseek-ai/dsh-llm'
  9. import { LlmAdapter } from '@deepseek-ai/dsh-llm'
  10. import type { SessionId } from '@deepseek-ai/dsh-session'
  11. import type { HostFrame, MuxFrame } from '@deepseek-ai/dsh-host-apiproxy/api'
  12. import type { RpcRequest, RpcResponse } from '@deepseek-ai/dsh-host-apiproxy/api/rpc'
  13. import { RpcId } from '@deepseek-ai/dsh-host-apiproxy/api/rpc'
  14. import { bootHost, startHost, type HostHandle, type RunningHost } from '../src/index.ts'
  15. /** Scripted adapter: each model call consumes the next chunk list; 'hang' streams then waits for abort. */
  16. class ScriptedAdapter extends LlmAdapter {
  17. constructor(private script: (StreamChunk[] | 'hang')[]) {
  18. super()
  19. }
  20. async * stream(options: GenerateOptions): AsyncIterable<StreamChunk> {
  21. const entry = this.script.shift()
  22. if (!entry) throw new Error('ScriptedAdapter: script exhausted')
  23. if (entry === 'hang') {
  24. yield { type: 'block-start', index: 0, blockType: 'text' }
  25. await new Promise<void>((_resolve, reject) => {
  26. options.signal?.addEventListener('abort', () => { reject(new Error('aborted')) }, { once: true })
  27. })
  28. return
  29. }
  30. yield * entry
  31. }
  32. }
  33. function textResponse(text: string): StreamChunk[] {
  34. return [
  35. { type: 'block-start', index: 0, blockType: 'text' },
  36. { type: 'text-delta', index: 0, text },
  37. { type: 'block-end', index: 0, block: { type: 'text', text } },
  38. { type: 'usage', usage: { inputTokens: 10, outputTokens: text.length } },
  39. { type: 'finish', reason: { kind: 'stop' } },
  40. ]
  41. }
  42. function request<P>(payload: P): RpcRequest<P> {
  43. return { rpcId: RpcId(`req-${String(nextRpc++)}`), payload }
  44. }
  45. let nextRpc = 1
  46. function waitForIdle(ctx: Context, agent: Agent): Promise<void> {
  47. return new Promise((resolve) => {
  48. const dispose = ctx.on('agent/status', (subject: Agent, status: string) => {
  49. if (subject === agent && status === 'idle') {
  50. dispose()
  51. resolve()
  52. }
  53. })
  54. })
  55. }
  56. function expectOk<T>(response: RpcResponse<T>): T {
  57. expect(response.result.ok).toBe(true)
  58. if (!response.result.ok) throw new Error('unreachable')
  59. return response.result.value
  60. }
  61. let host: RunningHost | undefined
  62. beforeEach(() => {
  63. vi.stubEnv('DEEPSEEK_API_KEY', 'spec-placeholder-key')
  64. })
  65. afterEach(async () => {
  66. await host?.dispose()
  67. host = undefined
  68. vi.unstubAllEnvs()
  69. })
  70. async function boot(script: (StreamChunk[] | 'hang')[] = []): Promise<RunningHost> {
  71. host = await startHost({
  72. boot: { persistenceRoot: mkdtempSync(join(tmpdir(), 'dsh-host-runtime-')), provider: 'scripted', model: 'test-model' },
  73. })
  74. host.ctx.llm.registerAdapter(['scripted'], new ScriptedAdapter(script))
  75. return host
  76. }
  77. describe('bootHost / startHost', () => {
  78. it('falls back to the deepseek defaults and disposes idempotently', async () => {
  79. const handle: HostHandle = await bootHost({ persistenceRoot: mkdtempSync(join(tmpdir(), 'dsh-boot-')) })
  80. expect(handle.defaults).toMatchObject({ provider: 'deepseek', model: 'deepseek-v4-flash' })
  81. expect(typeof handle.defaults.cwd).toBe('string')
  82. await handle.dispose()
  83. })
  84. it('startHost assembles api + handler over the same defaults and dedupes dispose', async () => {
  85. const running = await boot()
  86. expect(running.defaults).toMatchObject({ provider: 'scripted', model: 'test-model' })
  87. const body = JSON.stringify({ type: 'client-request', rpcId: 'r-h', method: 'host.describe', payload: {} })
  88. const response = await running.handler.fetch(new Request('http://x/api/host.describe', { method: 'POST', body }))
  89. const parsed = await response.json() as { result: { ok: boolean; value: { provider: string } } }
  90. expect(parsed.result.value.provider).toBe('scripted')
  91. const first = running.dispose()
  92. expect(running.dispose()).toBe(first)
  93. await first
  94. host = undefined
  95. })
  96. })
  97. describe('host.describe', () => {
  98. it('reports version, cwd, defaults, and the attached count', async () => {
  99. const { api } = await boot()
  100. const value = expectOk(await api.host.describe(request({})))
  101. expect(value).toMatchObject({ version: '0.0.1', cwd: process.cwd(), provider: 'scripted', model: 'test-model', attachedSessions: 0 })
  102. })
  103. })
  104. describe('sessions.create / list', () => {
  105. it('creates a session (echoing the request rpcId) and lists it newest-first', async () => {
  106. const { api } = await boot()
  107. const created = await api.sessions.create(request({ cwd: '/tmp' }))
  108. const { sessionId } = expectOk(created)
  109. expect(created.rpcId).toMatch(/^req-/)
  110. const second = expectOk(await api.sessions.create(request({}))).sessionId
  111. const { items } = expectOk(await api.sessions.list(request({})))
  112. expect(items.map(item => item.sessionId)).toContain(sessionId)
  113. expect(items.map(item => item.sessionId)).toContain(second)
  114. const first = items.find(item => item.sessionId === sessionId)
  115. expect(first?.cwd).toBe('/tmp')
  116. expect(first?.running).toBe(false)
  117. expect(first?.parentSessionId).toBeUndefined()
  118. })
  119. })
  120. describe('sessions.prompt / cancel', () => {
  121. it('queues a prompt whose rpcId rides into user/message, then the reply lands', async () => {
  122. const running = await boot([textResponse('pong')])
  123. const { api, ctx } = running
  124. const { sessionId } = expectOk(await api.sessions.create(request({})))
  125. const agent = ctx.agents.get(sessionId)
  126. expect(agent).toBeDefined()
  127. const idle = waitForIdle(ctx, agent as Agent)
  128. const promptRequest = request({ sessionId, mode: 'queue' as const, content: [{ type: 'text' as const, text: 'ping' }] })
  129. expectOk(await api.sessions.prompt(promptRequest))
  130. await idle
  131. const value = expectOk(await api.sessions.history(request({ sessionId })))
  132. const events = value.events.map(entry => entry.event)
  133. const userEvent = events.find(event => event.type === 'user/message') as
  134. | { data: { source?: { rpcId?: string } } } | undefined
  135. expect(userEvent?.data.source?.rpcId).toBe(promptRequest.rpcId)
  136. const reply = events.find(event => event.type === 'assistant/message')
  137. expect(reply).toBeDefined()
  138. })
  139. it('steer on an idle agent falls through to send', async () => {
  140. const running = await boot([textResponse('steered')])
  141. const { api, ctx } = running
  142. const { sessionId } = expectOk(await api.sessions.create(request({})))
  143. const idle = waitForIdle(ctx, ctx.agents.get(sessionId) as Agent)
  144. expectOk(await api.sessions.prompt(request({ sessionId, mode: 'steer' as const, content: [{ type: 'text' as const, text: 'now' }] })))
  145. await idle
  146. })
  147. it('errors session-not-found on a ghost session', async () => {
  148. const { api } = await boot()
  149. const response = await api.sessions.prompt(request({ sessionId: 'session-void' as SessionId, mode: 'queue' as const, content: [{ type: 'text' as const, text: 'x' }] }))
  150. expect(response.result.ok).toBe(false)
  151. if (!response.result.ok) expect(response.result.error.code).toBe('session-not-found')
  152. })
  153. it('maps a synchronous send throw to agent-busy', async () => {
  154. const { api } = await boot()
  155. const { sessionId } = expectOk(await api.sessions.create(request({})))
  156. const poisoned = [{ type: 'text', text: 'x', bad: () => 1 }] as never
  157. const response = await api.sessions.prompt(request({ sessionId, mode: 'queue' as const, content: poisoned }))
  158. expect(response.result.ok).toBe(false)
  159. if (!response.result.ok) expect(response.result.error.code).toBe('agent-busy')
  160. })
  161. it('cancels an attached agent and rejects an unattached one', async () => {
  162. const running = await boot(['hang'])
  163. const { api, ctx } = running
  164. const { sessionId } = expectOk(await api.sessions.create(request({})))
  165. const agent = ctx.agents.get(sessionId) as Agent
  166. agent.send([{ type: 'text', text: 'run forever' }])
  167. expectOk(await api.sessions.cancel(request({ sessionId })))
  168. const missing = await api.sessions.cancel(request({ sessionId: 'session-none' as SessionId }))
  169. expect(missing.result.ok).toBe(false)
  170. if (!missing.result.ok) expect(missing.result.error.code).toBe('session-not-found')
  171. })
  172. })
  173. describe('sessions.history', () => {
  174. it('implicitly resumes a cold session, deduplicating concurrent calls to one attach', async () => {
  175. const persistenceRoot = mkdtempSync(join(tmpdir(), 'dsh-host-resume-'))
  176. const first = await startHost({ boot: { persistenceRoot, provider: 'scripted', model: 'test-model' } })
  177. first.ctx.llm.registerAdapter(['scripted'], new ScriptedAdapter([textResponse('persisted')]))
  178. const { sessionId } = expectOk(await first.api.sessions.create(request({})))
  179. const agent = first.ctx.agents.get(sessionId) as Agent
  180. const idle = waitForIdle(first.ctx, agent)
  181. agent.send([{ type: 'text', text: 'save me' }])
  182. await idle
  183. await first.dispose()
  184. host = await startHost({ boot: { persistenceRoot, provider: 'scripted', model: 'test-model' } })
  185. host.ctx.llm.registerAdapter(['scripted'], new ScriptedAdapter([]))
  186. expect(host.ctx.agents.get(sessionId)).toBeUndefined()
  187. const [a, b] = await Promise.all([
  188. host.api.sessions.history(request({ sessionId })),
  189. host.api.sessions.history(request({ sessionId })),
  190. ])
  191. for (const response of [a, b]) {
  192. const value = expectOk(response)
  193. expect(value.events.some(entry => entry.event.type === 'assistant/message')).toBe(true)
  194. }
  195. expect(host.ctx.agents.get(sessionId)).toBeDefined()
  196. expect(host.ctx.agents.list()).toHaveLength(1)
  197. })
  198. it('errors session-not-found when resume fails, deduplicating concurrent resumes', async () => {
  199. const { api } = await boot()
  200. const ghost = 'session-ghost' as SessionId
  201. const [first, second] = await Promise.all([
  202. api.sessions.history(request({ sessionId: ghost })),
  203. api.sessions.history(request({ sessionId: ghost })),
  204. ])
  205. for (const response of [first, second]) {
  206. expect(response.result.ok).toBe(false)
  207. if (!response.result.ok) expect(response.result.error.code).toBe('session-not-found')
  208. }
  209. })
  210. it('paginates backwards on message boundaries with hasMore', async () => {
  211. const running = await boot([textResponse('a1'), textResponse('a2'), textResponse('a3')])
  212. const { api, ctx } = running
  213. const { sessionId } = expectOk(await api.sessions.create(request({})))
  214. const agent = ctx.agents.get(sessionId) as Agent
  215. for (const text of ['q1', 'q2', 'q3']) {
  216. const idle = waitForIdle(ctx, agent)
  217. agent.send([{ type: 'text', text }])
  218. await idle
  219. }
  220. const all = expectOk(await api.sessions.history(request({ sessionId })))
  221. expect(all.hasMore).toBe(false)
  222. const messageCount = all.events.filter(entry => entry.event.type === 'user/message' || entry.event.type === 'assistant/message').length
  223. expect(messageCount).toBe(6)
  224. const lastPage = expectOk(await api.sessions.history(request({ sessionId, maxMessages: 1 })))
  225. expect(lastPage.hasMore).toBe(true)
  226. expect(lastPage.events.filter(entry => entry.event.type === 'assistant/message')).toHaveLength(1)
  227. expect(lastPage.events.filter(entry => entry.event.type === 'user/message')).toHaveLength(0)
  228. const firstSeq = lastPage.events[0]?.event.seq as number
  229. const olderPage = expectOk(await api.sessions.history(request({ sessionId, beforeSeq: firstSeq, maxMessages: 2 })))
  230. expect(olderPage.events.at(-1)?.event.seq).toBeLessThan(firstSeq)
  231. expect(olderPage.hasMore).toBe(true)
  232. expect(olderPage.events.filter(entry => entry.event.type === 'user/message' || entry.event.type === 'assistant/message').length).toBe(2)
  233. })
  234. })
  235. describe('events streams', () => {
  236. it('mux: a pending pull wakes when a frame arrives (waiter path)', async () => {
  237. const running = await boot()
  238. const { api } = running
  239. const ac = new AbortController()
  240. const stream = api.events.mux(request({}), ac.signal)[Symbol.asyncIterator]()
  241. // no sessions yet: next() must pend on the queue's waiter, not the buffer
  242. const pending = stream.next()
  243. const { sessionId } = expectOk(await api.sessions.create(request({})))
  244. const frame = (await pending).value as RpcRequest<MuxFrame>
  245. expect(frame.payload).toMatchObject({ type: 'session/subscribed', sessionId })
  246. ac.abort()
  247. expect((await stream.next()).done).toBe(true)
  248. })
  249. it('lists fork lineage and announces it on the host stream', async () => {
  250. const running = await boot()
  251. const { api, ctx } = running
  252. const { sessionId: parent } = expectOk(await api.sessions.create(request({})))
  253. const ac = new AbortController()
  254. const stream = api.events.host(request({}), ac.signal)[Symbol.asyncIterator]()
  255. const child = `session-child-${String(Date.now())}` as SessionId
  256. const handle = await ctx.agents.create({ sessionId: child, meta: { parentSession: parent }, agentOptions: { provider: 'scripted', model: 'test-model' } })
  257. expect(handle.agent.id).toBe(child)
  258. const added = (await stream.next()).value as RpcRequest<HostFrame>
  259. expect(added.payload).toMatchObject({ type: 'host/session-added', sessionId: child, parentSessionId: parent })
  260. const { items } = expectOk(await api.sessions.list(request({})))
  261. expect(items.find(item => item.sessionId === child)?.parentSessionId).toBe(parent)
  262. await handle.dispose()
  263. let frame: RpcRequest<HostFrame>
  264. do frame = (await stream.next()).value as RpcRequest<HostFrame>
  265. while (frame.payload.type !== 'host/session-removed')
  266. expect(frame.payload).toMatchObject({ type: 'host/session-removed', sessionId: child })
  267. ac.abort()
  268. })
  269. it('mux: emits subscribed baselines, live session events, and new-session subscriptions until abort', async () => {
  270. const running = await boot([textResponse('live')])
  271. const { api, ctx } = running
  272. const { sessionId } = expectOk(await api.sessions.create(request({})))
  273. const ac = new AbortController()
  274. const stream = api.events.mux(request({}), ac.signal)[Symbol.asyncIterator]()
  275. const baseline = await stream.next()
  276. expect((baseline.value as RpcRequest<MuxFrame>).payload).toMatchObject({ type: 'session/subscribed', sessionId })
  277. const agent = ctx.agents.get(sessionId) as Agent
  278. const idle = waitForIdle(ctx, agent)
  279. agent.send([{ type: 'text', text: 'go' }])
  280. await idle
  281. const live = await stream.next()
  282. expect((live.value as RpcRequest<MuxFrame>).payload.type).toBe('session/event')
  283. const other = expectOk(await api.sessions.create(request({}))).sessionId
  284. let frame: RpcRequest<MuxFrame>
  285. do frame = (await stream.next()).value as RpcRequest<MuxFrame>
  286. while (!(frame.payload.type === 'session/subscribed' && frame.payload.sessionId === other))
  287. ac.abort()
  288. expect((await stream.next()).done).toBe(true)
  289. })
  290. it('host: session lifecycle, status flips (disposed suppressed), and agent errors', async () => {
  291. const running = await boot([textResponse('x')])
  292. const { api, ctx } = running
  293. const ac = new AbortController()
  294. const stream = api.events.host(request({}), ac.signal)[Symbol.asyncIterator]()
  295. const { sessionId } = expectOk(await api.sessions.create(request({})))
  296. const added = await stream.next()
  297. expect((added.value as RpcRequest<HostFrame>).payload).toMatchObject({ type: 'host/session-added', sessionId })
  298. const agent = ctx.agents.get(sessionId) as Agent
  299. const idle = waitForIdle(ctx, agent)
  300. agent.send([{ type: 'text', text: 'run' }])
  301. await idle
  302. const runningFrame = await stream.next()
  303. expect((runningFrame.value as RpcRequest<HostFrame>).payload).toMatchObject({ type: 'host/session-status', running: true })
  304. const idleFrame = await stream.next()
  305. expect((idleFrame.value as RpcRequest<HostFrame>).payload).toMatchObject({ type: 'host/session-status', running: false })
  306. // Raw ctx.emit lacks the scope carrier the mounted invariants plugin now
  307. // enforces; dispatch the way the loop does.
  308. agentEvents(ctx, agent).emit('agent/error', 1, 1, new Error('boom'))
  309. const errorFrame = await stream.next()
  310. expect((errorFrame.value as RpcRequest<HostFrame>).payload).toMatchObject({ type: 'host/agent-error', message: 'Error: boom' })
  311. ac.abort()
  312. // Push-after-done: an event landing between abort and generator wind-down
  313. // must be dropped silently, not crash the queue.
  314. agentEvents(ctx, agent).emit('agent/error', 1, 1, new Error('late'))
  315. expect((await stream.next()).done).toBe(true)
  316. })
  317. })
  318. describe('respond stub', () => {
  319. it('always reports not-pending (step2 registry pending)', async () => {
  320. const { api } = await boot()
  321. const receipt = await api.respond({ type: 'client-response', rpcId: RpcId('r'), result: { ok: true, value: null } })
  322. expect(receipt).toEqual({ accepted: false, reason: 'not-pending' })
  323. })
  324. })