host-runtime.spec.ts 35 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786
  1. import { existsSync, mkdirSync, mkdtempSync, writeFileSync } 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 { Config as SessionTitleConfig } from '@deepseek-ai/dsh-session-title'
  12. import type { Config as SessionTitleLlmConfig } from '@deepseek-ai/dsh-session-title-first-message-llm'
  13. import type { HostFrame, MuxFrame } from '@deepseek-ai/dsh-host-apiproxy/api'
  14. import type { RpcRequest, RpcResponse } from '@deepseek-ai/dsh-host-apiproxy/api/rpc'
  15. import { RpcId } from '@deepseek-ai/dsh-host-apiproxy/api/rpc'
  16. import { bootHost, startHost, type HostHandle, type RunningHost } from '../src/index.ts'
  17. /** Scripted adapter: each model call consumes the next chunk list; 'hang' streams then waits for abort. */
  18. class ScriptedAdapter extends LlmAdapter {
  19. readonly requests: GenerateOptions[] = []
  20. constructor(private script: (StreamChunk[] | 'hang')[]) {
  21. super()
  22. }
  23. async * stream(options: GenerateOptions): AsyncIterable<StreamChunk> {
  24. if ((options.tools?.length ?? 0) === 0) {
  25. yield * textResponse('Durable append-only session titles')
  26. return
  27. }
  28. this.requests.push(options)
  29. const entry = this.script.shift()
  30. if (!entry) throw new Error('ScriptedAdapter: script exhausted')
  31. if (entry === 'hang') {
  32. yield { type: 'block-start', index: 0, blockType: 'text' }
  33. await new Promise<void>((_resolve, reject) => {
  34. options.signal?.addEventListener('abort', () => { reject(new Error('aborted')) }, { once: true })
  35. })
  36. return
  37. }
  38. yield * entry
  39. }
  40. }
  41. function textResponse(text: string): StreamChunk[] {
  42. return [
  43. { type: 'block-start', index: 0, blockType: 'text' },
  44. { type: 'text-delta', index: 0, text },
  45. { type: 'block-end', index: 0, block: { type: 'text', text } },
  46. { type: 'usage', usage: { inputTokens: 10, outputTokens: text.length } },
  47. { type: 'finish', reason: { kind: 'stop' } },
  48. ]
  49. }
  50. function request<P>(payload: P): RpcRequest<P> {
  51. return { rpcId: RpcId(`req-${String(nextRpc++)}`), payload }
  52. }
  53. let nextRpc = 1
  54. function waitForIdle(ctx: Context, agent: Agent): Promise<void> {
  55. return new Promise((resolve) => {
  56. const dispose = ctx.on('agent/status', (subject: Agent, status: string) => {
  57. if (subject === agent && status === 'idle') {
  58. dispose()
  59. resolve()
  60. }
  61. })
  62. })
  63. }
  64. function expectOk<T>(response: RpcResponse<T>): T {
  65. expect(response.result.ok).toBe(true)
  66. if (!response.result.ok) throw new Error('unreachable')
  67. return response.result.value
  68. }
  69. async function nextMux(iterator: AsyncIterator<RpcRequest<MuxFrame>>): Promise<RpcRequest<MuxFrame>> {
  70. const next = await iterator.next()
  71. if (next.done === true) throw new Error('mux ended before the expected frame')
  72. return next.value
  73. }
  74. /** Durably append a title event without mounting title-generation policy. */
  75. function appendTitle(ctx: Context, agent: Agent, title: string) {
  76. return ctx.sessions.appendOutOfBand(agent.session, 'session/title', {
  77. title,
  78. messageSeqs: [1],
  79. source: { kind: 'fallback' },
  80. }, { kind: 'session-title' })
  81. }
  82. let host: RunningHost | undefined
  83. beforeEach(() => {
  84. vi.stubEnv('DEEPSEEK_API_KEY', 'spec-placeholder-key')
  85. })
  86. afterEach(async () => {
  87. await host?.dispose()
  88. host = undefined
  89. vi.unstubAllEnvs()
  90. })
  91. async function boot(
  92. script: (StreamChunk[] | 'hang')[] = [],
  93. sessionTitle?: SessionTitleConfig,
  94. sessionTitleLlm?: true | SessionTitleLlmConfig,
  95. ): Promise<RunningHost> {
  96. host = await startHost({
  97. boot: {
  98. persistenceRoot: mkdtempSync(join(tmpdir(), 'dsh-host-runtime-')),
  99. workspaceContext: false,
  100. provider: 'scripted',
  101. model: 'test-model',
  102. ...(sessionTitle === undefined ? {} : { sessionTitle }),
  103. ...(sessionTitleLlm === undefined ? {} : { sessionTitleLlm }),
  104. },
  105. })
  106. host.ctx.llm.registerAdapter(['scripted'], new ScriptedAdapter(script))
  107. return host
  108. }
  109. describe('bootHost / startHost', () => {
  110. it('falls back to the deepseek defaults and disposes idempotently', async () => {
  111. const handle: HostHandle = await bootHost({
  112. persistenceRoot: mkdtempSync(join(tmpdir(), 'dsh-boot-')),
  113. workspaceContext: false,
  114. })
  115. expect(handle.defaults).toMatchObject({ provider: 'deepseek', model: 'deepseek-v4-flash' })
  116. expect(typeof handle.defaults.cwd).toBe('string')
  117. await handle.dispose()
  118. })
  119. it('uses the JSONL backend compressed default', async () => {
  120. const handle: HostHandle = await bootHost({
  121. persistenceRoot: mkdtempSync(join(tmpdir(), 'dsh-boot-zstd-')),
  122. workspaceContext: false,
  123. })
  124. const session = handle.ctx.sessions.create()
  125. expect(handle.ctx.sessionPersistence.locate(session.header)?.path).toMatch(/\.jsonl\.zstd$/)
  126. await handle.dispose()
  127. })
  128. it('startHost assembles api + handler over the same defaults and dedupes dispose', async () => {
  129. const running = await boot()
  130. expect(running.defaults).toMatchObject({ provider: 'scripted', model: 'test-model' })
  131. const body = JSON.stringify({ type: 'client-request', rpcId: 'r-h', method: 'host.describe', payload: {} })
  132. const response = await running.handler.fetch(new Request('http://x/api/host.describe', { method: 'POST', body }))
  133. const parsed = await response.json() as { result: { ok: boolean; value: { provider: string } } }
  134. expect(parsed.result.value.provider).toBe('scripted')
  135. const first = running.dispose()
  136. expect(running.dispose()).toBe(first)
  137. await first
  138. host = undefined
  139. })
  140. it('routes workspace instructions through the assembled agent request prefix', async () => {
  141. const workspace = mkdtempSync(join(tmpdir(), 'dsh-host-workspace-'))
  142. mkdirSync(join(workspace, '.git'))
  143. writeFileSync(join(workspace, 'AGENTS.md'), 'host-workspace-context-probe\n')
  144. const adapter = new ScriptedAdapter([textResponse('done')])
  145. host = await startHost({
  146. boot: {
  147. persistenceRoot: mkdtempSync(join(tmpdir(), 'dsh-host-workspace-sessions-')),
  148. workspaceContext: { dshHome: join(workspace, '.dsh'), maxBytes: 65_536 },
  149. provider: 'scripted',
  150. model: 'test-model',
  151. cwd: workspace,
  152. },
  153. })
  154. host.ctx.llm.registerAdapter(['scripted'], adapter)
  155. const { sessionId } = expectOk(await host.api.sessions.create(request({})))
  156. const agent = host.ctx.agents.get(sessionId) as Agent
  157. const idle = waitForIdle(host.ctx, agent)
  158. expectOk(await host.api.sessions.prompt(request({
  159. sessionId,
  160. mode: 'queue' as const,
  161. content: [{ type: 'text' as const, text: 'go' }],
  162. })))
  163. await idle
  164. const requestText = adapter.requests[0]?.messages
  165. .flatMap(message => message.content)
  166. .filter(block => block.type === 'text')
  167. .map(block => block.text)
  168. .join('\n') ?? ''
  169. expect(requestText).toContain('Instructions from: AGENTS.md')
  170. expect(requestText).toContain('host-workspace-context-probe')
  171. })
  172. it('keeps model title generation disabled when sessionTitleLlm is omitted', async () => {
  173. const running = await boot([textResponse('pong')])
  174. const { api, ctx } = running
  175. const { sessionId } = expectOk(await api.sessions.create(request({})))
  176. const agent = ctx.agents.get(sessionId) as Agent
  177. const idle = waitForIdle(ctx, agent)
  178. expectOk(await api.sessions.prompt(request({
  179. sessionId,
  180. mode: 'queue' as const,
  181. content: [{ type: 'text' as const, text: 'Explain durable session titles.' }],
  182. })))
  183. await idle
  184. expect((await ctx.sessionTitle.refresh(agent.session))?.source).toEqual({ kind: 'fallback' })
  185. expect(agent.session.events.some(event => event.type === 'session/title-llm-request')).toBe(false)
  186. })
  187. })
  188. describe('host.describe', () => {
  189. it('reports version, cwd, defaults, and the attached count', async () => {
  190. const { api } = await boot()
  191. const value = expectOk(await api.host.describe(request({})))
  192. expect(value).toMatchObject({ version: '0.0.1', cwd: process.cwd(), provider: 'scripted', model: 'test-model', attachedSessions: 0 })
  193. })
  194. })
  195. describe('sessions.create / list', () => {
  196. it('creates a session (echoing the request rpcId) and lists it newest-first', async () => {
  197. const { api } = await boot()
  198. const created = await api.sessions.create(request({ cwd: '/tmp' }))
  199. const { sessionId } = expectOk(created)
  200. expect(created.rpcId).toMatch(/^req-/)
  201. const second = expectOk(await api.sessions.create(request({}))).sessionId
  202. const { items } = expectOk(await api.sessions.list(request({})))
  203. expect(items.map(item => item.sessionId)).toContain(sessionId)
  204. expect(items.map(item => item.sessionId)).toContain(second)
  205. const first = items.find(item => item.sessionId === sessionId)
  206. expect(first?.cwd).toBe('/tmp')
  207. expect(first?.running).toBe(false)
  208. expect(first?.parentSessionId).toBeUndefined()
  209. })
  210. it('ensures a missing project directory before minting the session', async () => {
  211. const { api } = await boot()
  212. const root = mkdtempSync(join(tmpdir(), 'dsh-host-create-cwd-'))
  213. const cwd = join(root, 'nested', 'workspace')
  214. expect(existsSync(cwd)).toBe(false)
  215. const { sessionId } = expectOk(await api.sessions.create(request({ cwd })))
  216. expect(existsSync(cwd)).toBe(true)
  217. const { items } = expectOk(await api.sessions.list(request({})))
  218. expect(items.find(item => item.sessionId === sessionId)?.cwd).toBe(cwd)
  219. })
  220. it('fails loud when the project directory cannot be created', async () => {
  221. const { api } = await boot()
  222. const root = mkdtempSync(join(tmpdir(), 'dsh-host-create-cwd-fail-'))
  223. const blocker = join(root, 'file-not-dir')
  224. writeFileSync(blocker, 'x')
  225. const response = await api.sessions.create(request({ cwd: join(blocker, 'child') }))
  226. expect(response.result.ok).toBe(false)
  227. if (response.result.ok) throw new Error('expected mkdir failure')
  228. expect(response.result.error.code).toBe('internal')
  229. expect(response.result.error.message).toMatch(/failed to ensure project directory/)
  230. })
  231. })
  232. describe('sessions.prompt / cancel', () => {
  233. it.each([
  234. { name: 'host default', config: true, target: '5 words', maxTokens: 64 },
  235. {
  236. name: 'configured policy',
  237. config: {
  238. targetWords: 3,
  239. targetCjkCharacters: 8,
  240. maxInputBytes: 2_048,
  241. maxOutputTokens: 24,
  242. timeoutMs: 2_000,
  243. },
  244. target: '3 words',
  245. maxTokens: 24,
  246. },
  247. ] satisfies {
  248. name: string
  249. config: true | SessionTitleLlmConfig
  250. target: string
  251. maxTokens: number
  252. }[])('replaces the fallback with a model-backed first-message title using the $name', async ({ config, target, maxTokens }) => {
  253. const modelTitle = 'Durable append-only session titles'
  254. const running = await boot([textResponse('pong')], undefined, config)
  255. const { api, ctx } = running
  256. const { sessionId } = expectOk(await api.sessions.create(request({})))
  257. const agent = ctx.agents.get(sessionId) as Agent
  258. const idle = waitForIdle(ctx, agent)
  259. expectOk(await api.sessions.prompt(request({
  260. sessionId,
  261. mode: 'queue' as const,
  262. content: [{ type: 'text' as const, text: 'Explain why append-only logs make session titles durable.' }],
  263. })))
  264. await idle
  265. await vi.waitFor(() => {
  266. expect(agent.session.events.filter(event => event.type === 'session/title').map(event => event.data))
  267. .toEqual([
  268. {
  269. title: 'Explain why append-only logs make',
  270. messageSeqs: [1],
  271. source: { kind: 'fallback' },
  272. },
  273. {
  274. title: modelTitle,
  275. messageSeqs: [1],
  276. source: {
  277. kind: 'provider',
  278. provider: 'session-title-first-message-llm',
  279. model: { provider: 'scripted', model: 'test-model' },
  280. },
  281. },
  282. ])
  283. })
  284. const titleRequest = agent.session.events.find(event => event.type === 'session/title-llm-request')
  285. expect(titleRequest?.data.system).toContain(target)
  286. expect(titleRequest?.data.maxTokens).toBe(maxTokens)
  287. })
  288. it.each([
  289. { name: 'host default', config: undefined, expected: 'Show the Web UI durable' },
  290. {
  291. name: 'configured limit',
  292. config: { fallbackMaxWords: 2, fallbackMaxBytes: 40, maxTitleBytes: 80 },
  293. expected: 'Show the',
  294. },
  295. ] satisfies { name: string; config: SessionTitleConfig | undefined; expected: string }[])(
  296. 'logs a durable fallback title with the $name',
  297. async ({ config, expected }) => {
  298. const running = await boot([textResponse('pong')], config)
  299. const { api, ctx } = running
  300. const { sessionId } = expectOk(await api.sessions.create(request({})))
  301. const agent = ctx.agents.get(sessionId) as Agent
  302. const idle = waitForIdle(ctx, agent)
  303. expectOk(await api.sessions.prompt(request({
  304. sessionId,
  305. mode: 'queue' as const,
  306. content: [{ type: 'text' as const, text: 'Show the Web UI durable session title' }],
  307. })))
  308. await idle
  309. const title = agent.session.events.find(event => event.type === 'session/title')
  310. expect(title?.data).toEqual({
  311. title: expected,
  312. messageSeqs: [1],
  313. source: { kind: 'fallback' },
  314. })
  315. },
  316. )
  317. it('queues a prompt whose rpcId rides into user/message, then the reply lands', async () => {
  318. const running = await boot([textResponse('pong')])
  319. const { api, ctx } = running
  320. const { sessionId } = expectOk(await api.sessions.create(request({})))
  321. const agent = ctx.agents.get(sessionId)
  322. expect(agent).toBeDefined()
  323. const idle = waitForIdle(ctx, agent as Agent)
  324. const promptRequest = request({ sessionId, mode: 'queue' as const, content: [{ type: 'text' as const, text: 'ping' }] })
  325. expectOk(await api.sessions.prompt(promptRequest))
  326. await idle
  327. const value = expectOk(await api.sessions.history(request({ sessionId })))
  328. const events = value.events.map(entry => entry.event)
  329. const userEvent = events.find(event => event.type === 'user/message') as
  330. | { data: { source?: { rpcId?: string } } } | undefined
  331. expect(userEvent?.data.source?.rpcId).toBe(promptRequest.rpcId)
  332. const reply = events.find(event => event.type === 'assistant/message')
  333. expect(reply).toBeDefined()
  334. })
  335. it('steer on an idle agent falls through to send', async () => {
  336. const running = await boot([textResponse('steered')])
  337. const { api, ctx } = running
  338. const { sessionId } = expectOk(await api.sessions.create(request({})))
  339. const idle = waitForIdle(ctx, ctx.agents.get(sessionId) as Agent)
  340. expectOk(await api.sessions.prompt(request({ sessionId, mode: 'steer' as const, content: [{ type: 'text' as const, text: 'now' }] })))
  341. await idle
  342. })
  343. it('errors session-not-found on a ghost session', async () => {
  344. const { api } = await boot()
  345. const response = await api.sessions.prompt(request({ sessionId: 'session-void' as SessionId, mode: 'queue' as const, content: [{ type: 'text' as const, text: 'x' }] }))
  346. expect(response.result.ok).toBe(false)
  347. if (!response.result.ok) expect(response.result.error.code).toBe('session-not-found')
  348. })
  349. it('cancels an attached agent and rejects an unattached one', async () => {
  350. const running = await boot(['hang'])
  351. const { api, ctx } = running
  352. const { sessionId } = expectOk(await api.sessions.create(request({})))
  353. const agent = ctx.agents.get(sessionId) as Agent
  354. agent.followup({ content: [{ type: 'text', text: 'run forever' }], source: { kind: 'user' } })
  355. expectOk(await api.sessions.cancel(request({ sessionId })))
  356. const missing = await api.sessions.cancel(request({ sessionId: 'session-none' as SessionId }))
  357. expect(missing.result.ok).toBe(false)
  358. if (!missing.result.ok) expect(missing.result.error.code).toBe('session-not-found')
  359. })
  360. })
  361. describe('sessions.history', () => {
  362. it('implicitly resumes a cold session, deduplicating concurrent calls to one attach', async () => {
  363. const persistenceRoot = mkdtempSync(join(tmpdir(), 'dsh-host-resume-'))
  364. const first = await startHost({
  365. boot: { persistenceRoot, workspaceContext: false, provider: 'scripted', model: 'test-model' },
  366. })
  367. first.ctx.llm.registerAdapter(['scripted'], new ScriptedAdapter([textResponse('persisted')]))
  368. const { sessionId } = expectOk(await first.api.sessions.create(request({})))
  369. const agent = first.ctx.agents.get(sessionId) as Agent
  370. const idle = waitForIdle(first.ctx, agent)
  371. agent.followup({ content: [{ type: 'text', text: 'save me' }], source: { kind: 'user' } })
  372. await idle
  373. const titleEvent = await appendTitle(first.ctx, agent, 'Persisted title')
  374. await first.dispose()
  375. host = await startHost({
  376. boot: { persistenceRoot, workspaceContext: false, provider: 'scripted', model: 'test-model' },
  377. })
  378. host.ctx.llm.registerAdapter(['scripted'], new ScriptedAdapter([]))
  379. expect(host.ctx.agents.get(sessionId)).toBeUndefined()
  380. const abort = new AbortController()
  381. const mux = host.api.events.mux(request({}), abort.signal)[Symbol.asyncIterator]()
  382. const [a, b] = await Promise.all([
  383. host.api.sessions.history(request({ sessionId })),
  384. host.api.sessions.history(request({ sessionId })),
  385. ])
  386. for (const response of [a, b]) {
  387. const value = expectOk(response)
  388. expect(value.events.some(entry => entry.event.type === 'assistant/message')).toBe(true)
  389. }
  390. expect(host.ctx.agents.get(sessionId)).toBeDefined()
  391. expect(host.ctx.agents.list()).toHaveLength(1)
  392. expect((await nextMux(mux)).payload).toMatchObject({ type: 'session/subscribed', sessionId })
  393. expect((await nextMux(mux)).payload).toEqual(expect.objectContaining({
  394. type: 'session/title', sessionId, title: 'Persisted title', eventSeq: titleEvent.seq,
  395. }))
  396. abort.abort()
  397. })
  398. it('errors session-not-found when resume fails, deduplicating concurrent resumes', async () => {
  399. const { api } = await boot()
  400. const ghost = 'session-ghost' as SessionId
  401. const [first, second] = await Promise.all([
  402. api.sessions.history(request({ sessionId: ghost })),
  403. api.sessions.history(request({ sessionId: ghost })),
  404. ])
  405. for (const response of [first, second]) {
  406. expect(response.result.ok).toBe(false)
  407. if (!response.result.ok) expect(response.result.error.code).toBe('session-not-found')
  408. }
  409. })
  410. it('paginates backwards on message boundaries with hasMore', async () => {
  411. const running = await boot([textResponse('a1'), textResponse('a2'), textResponse('a3')])
  412. const { api, ctx } = running
  413. const { sessionId } = expectOk(await api.sessions.create(request({})))
  414. const agent = ctx.agents.get(sessionId) as Agent
  415. for (const text of ['q1', 'q2', 'q3']) {
  416. const idle = waitForIdle(ctx, agent)
  417. agent.followup({ content: [{ type: 'text', text }], source: { kind: 'user' } })
  418. await idle
  419. }
  420. const all = expectOk(await api.sessions.history(request({ sessionId })))
  421. expect(all.hasMore).toBe(false)
  422. const messageCount = all.events.filter(entry =>
  423. entry.event.type === 'assistant/message'
  424. || (entry.event.type === 'user/message' && entry.event.data.source.kind === 'user'),
  425. ).length
  426. expect(messageCount).toBe(6)
  427. const lastPage = expectOk(await api.sessions.history(request({ sessionId, maxMessages: 1 })))
  428. expect(lastPage.hasMore).toBe(true)
  429. expect(lastPage.events.filter(entry => entry.event.type === 'assistant/message')).toHaveLength(1)
  430. expect(lastPage.events.filter(entry => entry.event.type === 'user/message')).toHaveLength(0)
  431. const firstSeq = lastPage.events[0]?.event.seq as number
  432. const olderPage = expectOk(await api.sessions.history(request({ sessionId, beforeSeq: firstSeq, maxMessages: 2 })))
  433. expect(olderPage.events.at(-1)?.event.seq).toBeLessThan(firstSeq)
  434. expect(olderPage.hasMore).toBe(true)
  435. expect(olderPage.events.filter(entry => entry.event.type === 'user/message' || entry.event.type === 'assistant/message').length).toBe(2)
  436. })
  437. })
  438. describe('events streams', () => {
  439. it('mux: a pending pull wakes when a frame arrives (waiter path)', async () => {
  440. const running = await boot()
  441. const { api } = running
  442. const ac = new AbortController()
  443. const stream = api.events.mux(request({}), ac.signal)[Symbol.asyncIterator]()
  444. // no sessions yet: next() must pend on the queue's waiter, not the buffer
  445. const pending = stream.next()
  446. const { sessionId } = expectOk(await api.sessions.create(request({})))
  447. const frame = (await pending).value as RpcRequest<MuxFrame>
  448. expect(frame.payload).toMatchObject({ type: 'session/subscribed', sessionId })
  449. ac.abort()
  450. expect((await stream.next()).done).toBe(true)
  451. })
  452. it('lists fork lineage and announces it on the host stream', async () => {
  453. const running = await boot()
  454. const { api, ctx } = running
  455. const { sessionId: parent } = expectOk(await api.sessions.create(request({})))
  456. const ac = new AbortController()
  457. const stream = api.events.host(request({}), ac.signal)[Symbol.asyncIterator]()
  458. const child = `session-child-${String(Date.now())}` as SessionId
  459. const handle = await ctx.agents.create({ sessionId: child, meta: { parentSession: parent }, agentOptions: { provider: 'scripted', model: 'test-model' } })
  460. expect(handle.agent.id).toBe(child)
  461. const added = (await stream.next()).value as RpcRequest<HostFrame>
  462. expect(added.payload).toMatchObject({ type: 'host/session-added', sessionId: child, parentSessionId: parent })
  463. const { items } = expectOk(await api.sessions.list(request({})))
  464. expect(items.find(item => item.sessionId === child)?.parentSessionId).toBe(parent)
  465. await handle.dispose()
  466. let frame: RpcRequest<HostFrame>
  467. do frame = (await stream.next()).value as RpcRequest<HostFrame>
  468. while (frame.payload.type !== 'host/session-removed')
  469. expect(frame.payload).toMatchObject({ type: 'host/session-removed', sessionId: child })
  470. ac.abort()
  471. })
  472. it('mux: emits subscribed baselines, live session events, and new-session subscriptions until abort', async () => {
  473. const running = await boot([textResponse('live')])
  474. const { api, ctx } = running
  475. const { sessionId } = expectOk(await api.sessions.create(request({})))
  476. const ac = new AbortController()
  477. const stream = api.events.mux(request({}), ac.signal)[Symbol.asyncIterator]()
  478. const baseline = await stream.next()
  479. expect((baseline.value as RpcRequest<MuxFrame>).payload).toMatchObject({ type: 'session/subscribed', sessionId })
  480. const agent = ctx.agents.get(sessionId) as Agent
  481. const idle = waitForIdle(ctx, agent)
  482. agent.followup({ content: [{ type: 'text', text: 'go' }], source: { kind: 'user' } })
  483. await idle
  484. const live = await stream.next()
  485. expect((live.value as RpcRequest<MuxFrame>).payload.type).toBe('session/event')
  486. const other = expectOk(await api.sessions.create(request({}))).sessionId
  487. let frame: RpcRequest<MuxFrame>
  488. do frame = (await stream.next()).value as RpcRequest<MuxFrame>
  489. while (!(frame.payload.type === 'session/subscribed' && frame.payload.sessionId === other))
  490. ac.abort()
  491. expect((await stream.next()).done).toBe(true)
  492. })
  493. it('mux: projects durable titles after open baselines and immediately after live raw events', async () => {
  494. const running = await boot()
  495. const { api, ctx } = running
  496. const { sessionId } = expectOk(await api.sessions.create(request({})))
  497. const agent = ctx.agents.get(sessionId) as Agent
  498. const initial = await appendTitle(ctx, agent, 'Initial title')
  499. const ac = new AbortController()
  500. const stream = api.events.mux(request({}), ac.signal)[Symbol.asyncIterator]()
  501. expect((await nextMux(stream)).payload).toMatchObject({ type: 'session/subscribed', sessionId })
  502. expect((await nextMux(stream)).payload).toEqual(expect.objectContaining({
  503. type: 'session/title', sessionId, title: 'Initial title', eventSeq: initial.seq, updatedAt: initial.time,
  504. }))
  505. const revised = await appendTitle(ctx, agent, 'Revised title')
  506. let raw: RpcRequest<MuxFrame>
  507. do raw = await nextMux(stream)
  508. while (!(raw.payload.type === 'session/event' && raw.payload.event.type === 'session/title'))
  509. expect(raw.payload).toMatchObject({ type: 'session/event', sessionId, event: { seq: revised.seq } })
  510. expect((await nextMux(stream)).payload).toEqual(expect.objectContaining({
  511. type: 'session/title', sessionId, title: 'Revised title', eventSeq: revised.seq, updatedAt: revised.time,
  512. }))
  513. ac.abort()
  514. })
  515. it('mux: emits no title control for untitled subscriptions', async () => {
  516. const { api } = await boot()
  517. const first = expectOk(await api.sessions.create(request({}))).sessionId
  518. const ac = new AbortController()
  519. const stream = api.events.mux(request({}), ac.signal)[Symbol.asyncIterator]()
  520. expect((await nextMux(stream)).payload).toMatchObject({ type: 'session/subscribed', sessionId: first })
  521. const second = expectOk(await api.sessions.create(request({}))).sessionId
  522. expect((await nextMux(stream)).payload).toMatchObject({ type: 'session/subscribed', sessionId: second })
  523. ac.abort()
  524. })
  525. it('host: session lifecycle, status flips (disposed suppressed), and agent errors', async () => {
  526. const running = await boot([textResponse('x')])
  527. const { api, ctx } = running
  528. const ac = new AbortController()
  529. const stream = api.events.host(request({}), ac.signal)[Symbol.asyncIterator]()
  530. const { sessionId } = expectOk(await api.sessions.create(request({})))
  531. const added = await stream.next()
  532. expect((added.value as RpcRequest<HostFrame>).payload).toMatchObject({ type: 'host/session-added', sessionId })
  533. const agent = ctx.agents.get(sessionId) as Agent
  534. const idle = waitForIdle(ctx, agent)
  535. agent.followup({ content: [{ type: 'text', text: 'run' }], source: { kind: 'user' } })
  536. await idle
  537. const runningFrame = await stream.next()
  538. expect((runningFrame.value as RpcRequest<HostFrame>).payload).toMatchObject({ type: 'host/session-status', running: true })
  539. const idleFrame = await stream.next()
  540. expect((idleFrame.value as RpcRequest<HostFrame>).payload).toMatchObject({ type: 'host/session-status', running: false })
  541. // Raw ctx.emit lacks the scope carrier the mounted invariants plugin now
  542. // enforces; dispatch the way the loop does.
  543. agentEvents(ctx, agent).emit('agent/error', 1, 1, 'boom')
  544. const errorFrame = await stream.next()
  545. expect((errorFrame.value as RpcRequest<HostFrame>).payload).toMatchObject({ type: 'host/agent-error', message: 'boom' })
  546. ac.abort()
  547. // Push-after-done: an event landing between abort and generator wind-down
  548. // must be dropped silently, not crash the queue.
  549. agentEvents(ctx, agent).emit('agent/error', 1, 1, new Error('late'))
  550. expect((await stream.next()).done).toBe(true)
  551. })
  552. })
  553. describe('question request / response', () => {
  554. const questions = [{
  555. id: 'mode', question: 'Choose a mode',
  556. options: [
  557. { label: 'Fast (Recommended)', description: 'Move quickly.' },
  558. { label: 'Careful', description: 'Review first.' },
  559. ],
  560. }]
  561. it('waits, replays the same rpcId on reconnect, validates, and resolves first-wins', async () => {
  562. const running = await boot()
  563. const { api, ctx } = running
  564. const { sessionId } = expectOk(await api.sessions.create(request({})))
  565. const agent = ctx.agents.get(sessionId) as Agent
  566. const ac = new AbortController()
  567. const stream = api.events.mux(request({}), ac.signal)[Symbol.asyncIterator]()
  568. await stream.next() // subscribed baseline starts the generator and installs the queue
  569. const answerPromise = ctx.userInteraction.ask({ questions, agent })
  570. const requested = (await stream.next()).value as RpcRequest<MuxFrame>
  571. expect(requested.payload).toMatchObject({ type: 'question/requested', sessionId, questions })
  572. const wrongSession = await api.respond({
  573. type: 'client-response', rpcId: requested.rpcId,
  574. result: {
  575. ok: true,
  576. value: { sessionId: 'session-other', answer: { answers: [{ id: 'mode', selected: ['Fast (Recommended)'] }] } },
  577. },
  578. })
  579. expect(wrongSession).toEqual({ accepted: false, reason: 'bad-response' })
  580. const badChoice = await api.respond({
  581. type: 'client-response', rpcId: requested.rpcId,
  582. result: {
  583. ok: true,
  584. value: { sessionId, answer: { answers: [{ id: 'mode', selected: ['Unknown'] }] } },
  585. },
  586. })
  587. expect(badChoice).toEqual({ accepted: false, reason: 'bad-response' })
  588. const invalidResults = [
  589. { ok: true as const, value: null },
  590. { ok: true as const, value: { sessionId, answer: { answers: [] } } },
  591. { ok: true as const, value: { sessionId, answer: { answers: [{ id: 'wrong', selected: ['Fast (Recommended)'] }] } } },
  592. { ok: true as const, value: { sessionId, answer: { answers: [{ id: 'mode', selected: ['Fast (Recommended)', 'Fast (Recommended)'] }] } } },
  593. { ok: true as const, value: { sessionId, answer: { answers: [{ id: 'mode', selected: ['Fast (Recommended)', 'Careful'] }] } } },
  594. { ok: true as const, value: { sessionId, answer: { answers: [{ id: 'mode', selected: [], custom: ' ' }] } } },
  595. { ok: true as const, value: { sessionId, answer: { answers: [{ id: 'mode', selected: ['Careful'], custom: 'Other' }] } } },
  596. { ok: false as const, error: { code: 'internal' as const, message: 'wrong error', details: {} } },
  597. ]
  598. for (const result of invalidResults) {
  599. expect(await api.respond({
  600. type: 'client-response', rpcId: requested.rpcId, result,
  601. })).toEqual({ accepted: false, reason: 'bad-response' })
  602. }
  603. const reconnectAbort = new AbortController()
  604. const replay = api.events.mux(request({}), reconnectAbort.signal)[Symbol.asyncIterator]()
  605. await replay.next()
  606. const replayed = (await replay.next()).value as RpcRequest<MuxFrame>
  607. expect(replayed.rpcId).toBe(requested.rpcId)
  608. expect(replayed.payload).toEqual(requested.payload)
  609. const response = {
  610. type: 'client-response' as const,
  611. rpcId: requested.rpcId,
  612. result: {
  613. ok: true as const,
  614. value: { sessionId, answer: { answers: [{ id: 'mode', selected: ['Fast (Recommended)'] }] } },
  615. },
  616. }
  617. const [first, duplicate] = await Promise.all([api.respond(response), api.respond(response)])
  618. expect([first, duplicate]).toContainEqual({ accepted: true })
  619. expect([first, duplicate]).toContainEqual({ accepted: false, reason: 'not-pending' })
  620. await expect(answerPromise).resolves.toEqual({
  621. answers: [{ id: 'mode', selected: ['Fast (Recommended)'] }],
  622. })
  623. const resolved = (await stream.next()).value as RpcRequest<MuxFrame>
  624. expect(resolved.payload).toMatchObject({
  625. type: 'question/resolved', sessionId, questionRpcId: requested.rpcId, outcome: 'answered',
  626. })
  627. expect(await api.respond(response)).toEqual({ accepted: false, reason: 'not-pending' })
  628. const customQuestions = [{ id: 'detail', question: 'What else?' }]
  629. const customAnswer = ctx.userInteraction.ask({ questions: customQuestions, agent })
  630. const customRequested = (await stream.next()).value as RpcRequest<MuxFrame>
  631. expect(await api.respond({
  632. type: 'client-response', rpcId: customRequested.rpcId,
  633. result: {
  634. ok: true,
  635. value: { sessionId, answer: { answers: [{ id: 'detail', selected: [], custom: 'Keep traces' }] } },
  636. },
  637. })).toEqual({ accepted: true })
  638. await expect(customAnswer).resolves.toEqual({
  639. answers: [{ id: 'detail', selected: [], custom: 'Keep traces' }],
  640. })
  641. expect(((await stream.next()).value as RpcRequest<MuxFrame>).payload).toMatchObject({
  642. type: 'question/resolved', questionRpcId: customRequested.rpcId, outcome: 'answered',
  643. })
  644. const blankAnswer = ctx.userInteraction.ask({ questions, agent })
  645. const blankRequested = (await stream.next()).value as RpcRequest<MuxFrame>
  646. expect(await api.respond({
  647. type: 'client-response', rpcId: blankRequested.rpcId,
  648. result: {
  649. ok: true,
  650. value: { sessionId, answer: { answers: [{ id: 'mode', selected: [] }] } },
  651. },
  652. })).toEqual({ accepted: true })
  653. await expect(blankAnswer).resolves.toEqual({
  654. answers: [{ id: 'mode', selected: [] }],
  655. })
  656. expect(((await stream.next()).value as RpcRequest<MuxFrame>).payload).toMatchObject({
  657. type: 'question/resolved', questionRpcId: blankRequested.rpcId, outcome: 'answered',
  658. })
  659. ac.abort()
  660. reconnectAbort.abort()
  661. })
  662. it('distinguishes user cancellation from owner abort and rejects late responses', async () => {
  663. const running = await boot()
  664. const { api, ctx } = running
  665. const { sessionId } = expectOk(await api.sessions.create(request({})))
  666. const agent = ctx.agents.get(sessionId) as Agent
  667. const streamAbort = new AbortController()
  668. const stream = api.events.mux(request({}), streamAbort.signal)[Symbol.asyncIterator]()
  669. await stream.next()
  670. const cancelled = ctx.userInteraction.ask({ questions, agent }).catch((error: unknown) => error)
  671. const requested = (await stream.next()).value as RpcRequest<MuxFrame>
  672. expect(await api.respond({
  673. type: 'client-response', rpcId: requested.rpcId,
  674. result: { ok: false, error: { code: 'cancelled', message: 'skip', details: {} } },
  675. })).toEqual({ accepted: true })
  676. await expect(cancelled).resolves.toMatchObject({ code: 'ASK_CANCELLED' })
  677. expect(((await stream.next()).value as RpcRequest<MuxFrame>).payload).toMatchObject({
  678. type: 'question/resolved', outcome: 'cancelled',
  679. })
  680. const ownerAbort = new AbortController()
  681. const aborted = ctx.userInteraction.ask({ questions, agent, signal: ownerAbort.signal })
  682. .catch((error: unknown) => error)
  683. const abortRequest = (await stream.next()).value as RpcRequest<MuxFrame>
  684. ownerAbort.abort()
  685. await expect(aborted).resolves.toMatchObject({ code: 'ASK_ABORTED' })
  686. expect(((await stream.next()).value as RpcRequest<MuxFrame>).payload).toMatchObject({
  687. type: 'question/resolved', questionRpcId: abortRequest.rpcId, outcome: 'cancelled',
  688. })
  689. expect(await api.respond({
  690. type: 'client-response', rpcId: abortRequest.rpcId,
  691. result: { ok: false, error: { code: 'cancelled', message: 'late', details: {} } },
  692. })).toEqual({ accepted: false, reason: 'not-pending' })
  693. streamAbort.abort()
  694. })
  695. it('rejects missing routing and pre-abort, then aborts outstanding waits on disposal', async () => {
  696. const running = await boot()
  697. const { ctx } = running
  698. await expect(ctx.userInteraction.ask({ questions })).rejects.toMatchObject({ code: 'ASK_MISSING_AGENT' })
  699. const { sessionId } = expectOk(await running.api.sessions.create(request({})))
  700. const agent = ctx.agents.get(sessionId) as Agent
  701. const alreadyAborted = new AbortController()
  702. alreadyAborted.abort()
  703. await expect(ctx.userInteraction.ask({ questions, agent, signal: alreadyAborted.signal }))
  704. .rejects.toMatchObject({ code: 'ASK_ABORTED' })
  705. const outstanding = ctx.userInteraction.ask({ questions, agent })
  706. const disposed = running.dispose()
  707. host = undefined
  708. await expect(outstanding).rejects.toMatchObject({ code: 'ASK_ABORTED' })
  709. await disposed
  710. })
  711. })