load.spec.ts 19 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337
  1. import { afterEach, beforeEach, describe, expect, it } from 'vitest'
  2. import { mkdtemp, rm } from 'node:fs/promises'
  3. import { tmpdir } from 'node:os'
  4. import { join } from 'node:path'
  5. import { PROTOCOL_VERSION } from '@agentclientprotocol/sdk'
  6. import { SESSION_FORMAT_VERSION, SessionId } from '@deepseek-ai/dsh-session'
  7. import type {} from '@deepseek-ai/dsh-session-title'
  8. import { makeBridgeHarness, textResponse, toolCallResponse, type BridgeHarness, type CapturedUpdate } from './harness.ts'
  9. /** Concatenate the text of all agent_message_chunk updates. */
  10. function messageText(updates: CapturedUpdate[]): string {
  11. return updates
  12. .filter(u => u.sessionUpdate === 'agent_message_chunk')
  13. .map(u => (u.content.type === 'text' ? u.content.text : ''))
  14. .join('')
  15. }
  16. describe('acp bridge — session/load replay', () => {
  17. let storageDir: string
  18. let live: BridgeHarness | undefined
  19. let loader: BridgeHarness | undefined
  20. beforeEach(async () => { storageDir = await mkdtemp(join(tmpdir(), 'acp-load-')) })
  21. afterEach(async () => {
  22. if (live) await live.dispose()
  23. if (loader) await loader.dispose()
  24. live = loader = undefined
  25. await rm(storageDir, { recursive: true, force: true })
  26. })
  27. it('replays a persisted turn from the event log as session/update on load', async () => {
  28. // 1. Create a session and run one turn — persistence writes the event log.
  29. live = await makeBridgeHarness({ storageDir, script: [textResponse('remembered answer')] })
  30. await live.client.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: {} })
  31. const { sessionId } = await live.client.newSession({ cwd: process.cwd(), mcpServers: [] })
  32. await live.client.prompt({ sessionId, prompt: [{ type: 'text', text: 'remember this' }] })
  33. // Dispose to flush + release; the on-disk log persists.
  34. await live.dispose()
  35. live = undefined
  36. // 2. A fresh bridge loads the same session id and must replay the turn.
  37. loader = await makeBridgeHarness({ storageDir, script: [] })
  38. await loader.client.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: {} })
  39. const res = await loader.client.loadSession({ sessionId, cwd: process.cwd(), mcpServers: [] })
  40. expect(res).toBeDefined()
  41. // The replayed updates reconstruct the assistant text from the event log
  42. // (assistant/chunk → agent_message_chunk), NOT from deriveMessages.
  43. expect(messageText(loader.updates)).toBe('remembered answer')
  44. // And the USER side of the turn replays too (user/message →
  45. // user_message_chunk), so the editor transcript shows both sides.
  46. const userText = loader.updates
  47. .filter(u => u.sessionUpdate === 'user_message_chunk')
  48. .map(u => (u.content.type === 'text' ? u.content.text : ''))
  49. .join('')
  50. expect(userText).toBe('remember this')
  51. })
  52. it('streams and replays the same persisted session_info_update for a title event', async () => {
  53. live = await makeBridgeHarness({ storageDir, script: [] })
  54. await live.client.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: {} })
  55. const { sessionId } = await live.client.newSession({ cwd: process.cwd(), mcpServers: [] })
  56. const session = live.ctx.agents.get(SessionId(sessionId))!.session
  57. const event = await live.ctx.sessions.appendOutOfBand(session, 'session/title', {
  58. title: 'Durable ACP title',
  59. messageSeqs: [1],
  60. source: { kind: 'fallback' },
  61. }, { kind: 'session-title' })
  62. const expected = {
  63. sessionUpdate: 'session_info_update' as const,
  64. title: 'Durable ACP title',
  65. updatedAt: new Date(event.time).toISOString(),
  66. }
  67. expect(live.updates).toContainEqual(expected)
  68. await live.dispose()
  69. live = undefined
  70. loader = await makeBridgeHarness({ storageDir, script: [] })
  71. await loader.client.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: {} })
  72. await loader.client.loadSession({ sessionId, cwd: process.cwd(), mcpServers: [] })
  73. expect(loader.updates).toContainEqual(expected)
  74. })
  75. it('replays a persisted tool call with the TOOL-OWNED presentation (title/rawInput/console output)', async () => {
  76. // Persist a real bash call, then replay it through a fresh bridge. A throwaway presenter pairs
  77. // call and result in log order so replay uses the shipping tool's same cards as live streaming.
  78. live = await makeBridgeHarness({
  79. storageDir,
  80. withBash: true,
  81. script: [toolCallResponse('c1', 'bash', { command: 'echo hello', description: 'Print a greeting' }), textResponse('done')],
  82. })
  83. await live.client.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: {} })
  84. const { sessionId } = await live.client.newSession({ cwd: process.cwd(), mcpServers: [] })
  85. await live.client.prompt({ sessionId, prompt: [{ type: 'text', text: 'greet' }] })
  86. await live.dispose()
  87. live = undefined
  88. // A fresh bridge — also with the real bash tool, since the presentation is
  89. // resolved from the live registry at replay time — loads the session.
  90. loader = await makeBridgeHarness({ storageDir, withBash: true, script: [] })
  91. await loader.client.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: {} })
  92. await loader.client.loadSession({ sessionId, cwd: process.cwd(), mcpServers: [] })
  93. const call = loader.updates.find(u => u.sessionUpdate === 'tool_call')
  94. expect(call).toMatchObject({ toolCallId: 'c1', title: 'echo hello', kind: 'execute', rawInput: 'echo hello' })
  95. if (call?.sessionUpdate !== 'tool_call') throw new Error('expected a tool_call')
  96. // Capability OFF on this loader: the description renders as a content block, no terminal block.
  97. expect(call.content).toEqual([{ type: 'content', content: { type: 'text', text: 'Print a greeting' } }])
  98. const update = loader.updates.find(u => u.sessionUpdate === 'tool_call_update')
  99. expect(update?.sessionUpdate).toBe('tool_call_update')
  100. if (update?.sessionUpdate !== 'tool_call_update') throw new Error('expected a tool_call_update')
  101. expect(update).toMatchObject({ toolCallId: 'c1', status: 'completed' })
  102. const content = update.content as { content: { text: string } }[]
  103. expect(content[0]?.content.text).toBe('```console\nhello\n```')
  104. })
  105. it('replays a persisted todo/write as a plan sessionUpdate on load', async () => {
  106. // A persisted `todo/write` must replay as an ACP plan update so a reopened editor sees the
  107. // current plan, not just the tool transcript.
  108. live = await makeBridgeHarness({
  109. storageDir,
  110. withTodo: true,
  111. script: [
  112. toolCallResponse('c1', 'todo_write', {
  113. todos: [
  114. { content: 'first step', status: 'in_progress' },
  115. { content: 'second step', status: 'pending' },
  116. ],
  117. }),
  118. textResponse('planned'),
  119. ],
  120. })
  121. await live.client.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: {} })
  122. const { sessionId } = await live.client.newSession({ cwd: process.cwd(), mcpServers: [] })
  123. await live.client.prompt({ sessionId, prompt: [{ type: 'text', text: 'plan it' }] })
  124. await live.dispose()
  125. live = undefined
  126. loader = await makeBridgeHarness({ storageDir, withTodo: true, script: [] })
  127. await loader.client.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: {} })
  128. await loader.client.loadSession({ sessionId, cwd: process.cwd(), mcpServers: [] })
  129. const plan = loader.updates.find(u => u.sessionUpdate === 'plan')
  130. expect(plan).toEqual({
  131. sessionUpdate: 'plan',
  132. entries: [
  133. { content: 'first step', priority: 'medium', status: 'in_progress' },
  134. { content: 'second step', priority: 'medium', status: 'pending' },
  135. ],
  136. })
  137. })
  138. it('replays a persisted bash call as a TERMINAL card when the loader advertises the capability', async () => {
  139. // The presentation is resolved at replay time, so a loader that advertised
  140. // _meta.terminal_output must reconstruct the terminal card (content + _meta)
  141. // from the persisted log — identical to how it would have streamed live.
  142. live = await makeBridgeHarness({
  143. storageDir,
  144. withBash: true,
  145. script: [toolCallResponse('c1', 'bash', { command: 'echo hi', description: 'Greet' }), textResponse('done')],
  146. })
  147. await live.client.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: {} })
  148. const { sessionId } = await live.client.newSession({ cwd: process.cwd(), mcpServers: [] })
  149. await live.client.prompt({ sessionId, prompt: [{ type: 'text', text: 'greet' }] })
  150. await live.dispose()
  151. live = undefined
  152. loader = await makeBridgeHarness({ storageDir, withBash: true, script: [] })
  153. await loader.client.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: { _meta: { terminal_output: true } } })
  154. await loader.client.loadSession({ sessionId, cwd: process.cwd(), mcpServers: [] })
  155. const call = loader.updates.find(u => u.sessionUpdate === 'tool_call')
  156. if (call?.sessionUpdate !== 'tool_call') throw new Error('expected a tool_call')
  157. // Replay reconstructs the terminal card: description block, then terminal block.
  158. expect(call.content).toEqual([
  159. { type: 'content', content: { type: 'text', text: 'Greet' } },
  160. { type: 'terminal', terminalId: 'c1' },
  161. ])
  162. expect((call._meta as { terminal_info?: unknown }).terminal_info).toEqual({ terminal_id: 'c1', cwd: process.cwd() })
  163. const update = loader.updates.find(u => u.sessionUpdate === 'tool_call_update')
  164. if (update?.sessionUpdate !== 'tool_call_update') throw new Error('expected a tool_call_update')
  165. // Terminal mode: content omitted, output + exit on _meta — matching live.
  166. expect(update.content).toBeUndefined()
  167. const meta = update._meta as { terminal_output?: { data: string }; terminal_exit?: { exit_code?: number } }
  168. expect(meta.terminal_output?.data).toBe('hi\n')
  169. expect(meta.terminal_exit?.exit_code).toBe(0)
  170. })
  171. it('keeps one terminal completion live and on replay when a pruning replacement is logged', async () => {
  172. live = await makeBridgeHarness({
  173. storageDir,
  174. withBash: true,
  175. script: [toolCallResponse('c1', 'bash', { command: 'echo full', description: 'Print full output' }), textResponse('done')],
  176. })
  177. await live.client.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: { _meta: { terminal_output: true } } })
  178. const { sessionId } = await live.client.newSession({ cwd: process.cwd(), mcpServers: [] })
  179. await live.client.prompt({ sessionId, prompt: [{ type: 'text', text: 'run it' }] })
  180. const session = live.ctx.agents.get(SessionId(sessionId))!.session
  181. const original = session.events.find(event => event.type === 'tool/result')
  182. if (original?.type !== 'tool/result') throw new Error('expected original tool/result')
  183. const liveCompletions = () => live!.updates.filter(update =>
  184. update.sessionUpdate === 'tool_call_update' && update.toolCallId === 'c1')
  185. expect(liveCompletions()).toHaveLength(1)
  186. expect((liveCompletions()[0] as { _meta?: { terminal_output?: { data: string } } })._meta?.terminal_output?.data)
  187. .toBe('full\n')
  188. session.append('turn/start', { turn: 2, trigger: { kind: 'message', source: { kind: 'user' } } })
  189. session.append('tool/result', {
  190. ...original.data,
  191. content: [{ type: 'text', text: '[... tool result middle pruned ...]' }],
  192. }, {
  193. surfaceOp: { op: 'replace', start: original.seq, end: original.seq },
  194. sourceEventSeqs: [original.seq],
  195. })
  196. session.append('turn/end', { turn: 2, reason: { kind: 'completed' } })
  197. // The replacement is durable but is not another live completion.
  198. expect(session.events.filter(event => event.type === 'tool/result')).toHaveLength(2)
  199. expect(JSON.stringify(session.deriveMessages())).toContain('tool result middle pruned')
  200. expect(liveCompletions()).toHaveLength(1)
  201. await live.dispose()
  202. live = undefined
  203. loader = await makeBridgeHarness({ storageDir, withBash: true, script: [] })
  204. await loader.client.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: { _meta: { terminal_output: true } } })
  205. await loader.client.loadSession({ sessionId, cwd: process.cwd(), mcpServers: [] })
  206. const replayed = loader.updates.filter(update =>
  207. update.sessionUpdate === 'tool_call_update' && update.toolCallId === 'c1')
  208. expect(replayed).toHaveLength(1)
  209. expect((replayed[0] as { _meta?: { terminal_output?: { data: string } } })._meta?.terminal_output?.data)
  210. .toBe('full\n')
  211. })
  212. it('a load whose resume finishes after a client disconnect leaks no live session', async () => {
  213. // Stall persistence so transport closes while resume is pending. Whether the SDK rejects first
  214. // or the bridge's post-await guard fires, no agent may survive for the dead connection.
  215. live = await makeBridgeHarness({ storageDir, script: [textResponse('x')] })
  216. await live.client.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: {} })
  217. const { sessionId } = await live.client.newSession({ cwd: process.cwd(), mcpServers: [] })
  218. await live.client.prompt({ sessionId, prompt: [{ type: 'text', text: 'hi' }] })
  219. await live.dispose()
  220. live = undefined
  221. loader = await makeBridgeHarness({ storageDir, script: [] })
  222. await loader.client.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: {} })
  223. const realLoad = loader.ctx.sessionPersistence.load.bind(loader.ctx.sessionPersistence)
  224. let release!: () => void
  225. const gate = new Promise<void>((r) => { release = r })
  226. loader.ctx.sessionPersistence.load = async (id) => { await gate; return realLoad(id) }
  227. const loadResult = loader.client.loadSession({ sessionId, cwd: process.cwd(), mcpServers: [] })
  228. .then(() => 'resolved' as const, () => 'rejected' as const)
  229. await loader.closeClientTransport() // teardown sets `closed` while load is gated
  230. release() // resume() finishes AFTER teardown
  231. expect(await loadResult).toBe('rejected')
  232. // No live agent was installed for the closed connection.
  233. expect(loader.ctx.agents.get(SessionId(sessionId))).toBeUndefined()
  234. })
  235. it('rejects load when the requested cwd does not match the persisted session cwd', async () => {
  236. // Seed a session on disk whose header.cwd is a DIFFERENT absolute path than the server's
  237. // launch dir. Resume must retain the header cwd and route bash there rather than reject the
  238. // mismatch or substitute the server cwd.
  239. loader = await makeBridgeHarness({ storageDir, script: [] })
  240. const otherCwd = '/some/other/workspace'
  241. await loader.ctx.sessionPersistence.create({
  242. version: SESSION_FORMAT_VERSION, id: SessionId('elsewhere'), createdAt: 1, cwd: otherCwd,
  243. })
  244. await loader.ctx.sessionPersistence.append(SessionId('elsewhere'), [
  245. { type: 'turn/start', seq: 0, time: 0, data: { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } } },
  246. { type: 'turn/end', seq: 1, time: 0, data: { turn: 1, reason: { kind: 'completed' } } },
  247. ])
  248. await loader.client.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: {} })
  249. await expect(loader.client.loadSession({ sessionId: 'elsewhere', cwd: process.cwd(), mcpServers: [] }))
  250. .rejects.toThrow(/cwd mismatch/)
  251. expect(loader.ctx.agents.get(SessionId('elsewhere'))).toBeUndefined()
  252. const res = await loader.client.loadSession({ sessionId: 'elsewhere', cwd: `${otherCwd}/.`, mcpServers: [] })
  253. expect(res).toBeDefined()
  254. expect(loader.ctx.agents.get(SessionId('elsewhere'))!.session.header.cwd).toBe(otherCwd)
  255. })
  256. it('rejects load for a non-absolute cwd (still required to be absolute)', async () => {
  257. loader = await makeBridgeHarness({ storageDir, script: [] })
  258. await loader.client.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: {} })
  259. await expect(loader.client.loadSession({ sessionId: 's', cwd: 'rel', mcpServers: [] }))
  260. .rejects.toThrow(/absolute/)
  261. })
  262. it('lets persistence reject a load for an unknown id after metadata lookup misses', async () => {
  263. loader = await makeBridgeHarness({ storageDir, script: [] })
  264. await loader.client.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: {} })
  265. await expect(loader.client.loadSession({ sessionId: 'missing', cwd: process.cwd(), mcpServers: [] }))
  266. .rejects.toThrow(/Internal error/)
  267. })
  268. it('rejects loading a persisted session that has NO cwd (would silently run in the launch dir)', async () => {
  269. // A legacy/external log without `header.cwd` must be rejected; the request cwd does not override
  270. // it, and accepting would let bash silently fall back to the server launch directory.
  271. loader = await makeBridgeHarness({ storageDir, script: [] })
  272. await loader.ctx.sessionPersistence.create({
  273. version: SESSION_FORMAT_VERSION, id: SessionId('legacy'), createdAt: 1, // no cwd
  274. })
  275. await loader.ctx.sessionPersistence.append(SessionId('legacy'), [
  276. { type: 'turn/start', seq: 0, time: 0, data: { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } } },
  277. { type: 'turn/end', seq: 1, time: 0, data: { turn: 1, reason: { kind: 'completed' } } },
  278. ])
  279. await loader.client.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: {} })
  280. await expect(loader.client.loadSession({ sessionId: 'legacy', cwd: process.cwd(), mcpServers: [] }))
  281. .rejects.toThrow(/no absolute persisted cwd/)
  282. // Rejected BEFORE resume (metadata-only check) — no agent was registered, so
  283. // the id is not wedged: a later attempt hits the same clean rejection, not a
  284. // duplicate-registration error.
  285. expect(loader.ctx.agents.get(SessionId('legacy'))).toBeUndefined()
  286. await expect(loader.client.loadSession({ sessionId: 'legacy', cwd: process.cwd(), mcpServers: [] }))
  287. .rejects.toThrow(/no absolute persisted cwd/)
  288. })
  289. it('allows loading alongside an existing session but rejects re-loading the SAME id', async () => {
  290. // Multi-session: a load can coexist with a live session, but loading an id
  291. // that is already live is rejected (it is already loaded).
  292. live = await makeBridgeHarness({ storageDir, script: [textResponse('one')] })
  293. await live.client.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: {} })
  294. const { sessionId } = await live.client.newSession({ cwd: process.cwd(), mcpServers: [] })
  295. await live.client.prompt({ sessionId, prompt: [{ type: 'text', text: 'hi' }] })
  296. // A different new session coexists.
  297. const other = await live.client.newSession({ cwd: process.cwd(), mcpServers: [] })
  298. expect(other.sessionId).not.toBe(sessionId)
  299. // Re-loading the already-live id is rejected.
  300. await expect(live.client.loadSession({ sessionId, cwd: process.cwd(), mcpServers: [] }))
  301. .rejects.toThrow(/already loaded/)
  302. })
  303. })