multi-session.spec.ts 7.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129
  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 { makeBridgeHarness, textResponse, type BridgeHarness, type CapturedUpdate } from './harness.ts'
  7. import { SessionId } from '@deepseek-ai/dsh-session'
  8. /** Text of the agent_message_chunk updates scoped to one session id. */
  9. function messageTextFor(updates: { sessionId?: string; update: CapturedUpdate }[], sessionId: string): string {
  10. return updates
  11. .filter(u => u.sessionId === sessionId && u.update.sessionUpdate === 'agent_message_chunk')
  12. .map(u => (u.update.sessionUpdate === 'agent_message_chunk' && u.update.content.type === 'text' ? u.update.content.text : ''))
  13. .join('')
  14. }
  15. describe('acp bridge — multi-session isolation', () => {
  16. let storageDir: string
  17. let harness: BridgeHarness | undefined
  18. beforeEach(async () => { storageDir = await mkdtemp(join(tmpdir(), 'acp-multi-')) })
  19. afterEach(async () => {
  20. if (harness) await harness.dispose()
  21. harness = undefined
  22. await rm(storageDir, { recursive: true, force: true })
  23. })
  24. it('two sessions stream concurrently without interleaving their updates', async () => {
  25. // Each session's prompt answer must arrive only on its own sessionId. The
  26. // scripted adapter answers in send order; both prompts run, and the bridge
  27. // demuxes every chunk by session id.
  28. harness = await makeBridgeHarness({ storageDir, script: [textResponse('answer-A'), textResponse('answer-B')] })
  29. await harness.client.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: {} })
  30. const a = (await harness.client.newSession({ cwd: process.cwd(), mcpServers: [] })).sessionId
  31. const b = (await harness.client.newSession({ cwd: process.cwd(), mcpServers: [] })).sessionId
  32. const [ra, rb] = await Promise.all([
  33. harness.client.prompt({ sessionId: a, prompt: [{ type: 'text', text: 'go A' }] }),
  34. harness.client.prompt({ sessionId: b, prompt: [{ type: 'text', text: 'go B' }] }),
  35. ])
  36. expect(ra.stopReason).toBe('end_turn')
  37. expect(rb.stopReason).toBe('end_turn')
  38. // A's text landed only on A; B's only on B (strict id demux, no interleave).
  39. expect(messageTextFor(harness.sessionUpdates, a)).toContain('answer-A')
  40. expect(messageTextFor(harness.sessionUpdates, a)).not.toContain('answer-B')
  41. expect(messageTextFor(harness.sessionUpdates, b)).toContain('answer-B')
  42. expect(messageTextFor(harness.sessionUpdates, b)).not.toContain('answer-A')
  43. })
  44. it('cancel in one session leaves the other session untouched', async () => {
  45. // Session A hangs; session B completes normally. Cancelling A settles ONLY
  46. // A as cancelled and never disturbs B's stream or result.
  47. harness = await makeBridgeHarness({ storageDir, script: ['hang', textResponse('B done')] })
  48. await harness.client.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: {} })
  49. const a = (await harness.client.newSession({ cwd: process.cwd(), mcpServers: [] })).sessionId
  50. const b = (await harness.client.newSession({ cwd: process.cwd(), mcpServers: [] })).sessionId
  51. const aPromise = harness.client.prompt({ sessionId: a, prompt: [{ type: 'text', text: 'hang A' }] })
  52. await new Promise(r => setTimeout(r, 30))
  53. await harness.client.cancel({ sessionId: a })
  54. expect((await aPromise).stopReason).toBe('cancelled')
  55. // B runs to completion, unaffected by A's cancel.
  56. const rb = await harness.client.prompt({ sessionId: b, prompt: [{ type: 'text', text: 'go B' }] })
  57. expect(rb.stopReason).toBe('end_turn')
  58. expect(messageTextFor(harness.sessionUpdates, b)).toContain('B done')
  59. })
  60. it('enforces one in-flight prompt PER session independently', async () => {
  61. harness = await makeBridgeHarness({ storageDir, script: ['hang', 'hang'] })
  62. await harness.client.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: {} })
  63. const a = (await harness.client.newSession({ cwd: process.cwd(), mcpServers: [] })).sessionId
  64. const b = (await harness.client.newSession({ cwd: process.cwd(), mcpServers: [] })).sessionId
  65. // One in-flight prompt in EACH session is allowed (independent limits).
  66. const aPromise = harness.client.prompt({ sessionId: a, prompt: [{ type: 'text', text: 'one A' }] })
  67. const bPromise = harness.client.prompt({ sessionId: b, prompt: [{ type: 'text', text: 'one B' }] })
  68. await new Promise(r => setTimeout(r, 30))
  69. // A second prompt in A is rejected, but B's in-flight prompt is unaffected.
  70. await expect(harness.client.prompt({ sessionId: a, prompt: [{ type: 'text', text: 'two A' }] }))
  71. .rejects.toThrow(/already in flight/)
  72. await harness.client.cancel({ sessionId: a })
  73. await harness.client.cancel({ sessionId: b })
  74. expect((await aPromise).stopReason).toBe('cancelled')
  75. expect((await bPromise).stopReason).toBe('cancelled')
  76. })
  77. it('a cancel for a non-existent session id is a silent no-op (does not touch others)', async () => {
  78. harness = await makeBridgeHarness({ storageDir, script: [textResponse('A done')] })
  79. await harness.client.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: {} })
  80. const a = (await harness.client.newSession({ cwd: process.cwd(), mcpServers: [] })).sessionId
  81. await expect(harness.client.cancel({ sessionId: 'ghost' })).resolves.toBeUndefined()
  82. // A still works after a cancel for an unknown id.
  83. const ra = await harness.client.prompt({ sessionId: a, prompt: [{ type: 'text', text: 'go A' }] })
  84. expect(ra.stopReason).toBe('end_turn')
  85. })
  86. it('disposing the whole bridge drains all live sessions to quiescence', async () => {
  87. harness = await makeBridgeHarness({ storageDir, script: ['hang', 'hang'] })
  88. await harness.client.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: {} })
  89. const a = (await harness.client.newSession({ cwd: process.cwd(), mcpServers: [] })).sessionId
  90. const b = (await harness.client.newSession({ cwd: process.cwd(), mcpServers: [] })).sessionId
  91. const agentA = harness.ctx.agents.get(SessionId(a))!
  92. const agentB = harness.ctx.agents.get(SessionId(b))!
  93. // Wait deterministically for BOTH agents to enter `running` (not a fixed
  94. // sleep — agent startup latency is unbounded on a loaded worker).
  95. const running = (agent: typeof agentA) => agent.status === 'running'
  96. ? Promise.resolve()
  97. : new Promise<void>((resolve) => {
  98. const dispose = harness!.ctx.on('agent/status', (subject, status) => {
  99. if (subject === agent && status === 'running') { dispose(); resolve() }
  100. })
  101. })
  102. void harness.client.prompt({ sessionId: a, prompt: [{ type: 'text', text: 'go A' }] }).catch(() => {})
  103. void harness.client.prompt({ sessionId: b, prompt: [{ type: 'text', text: 'go B' }] }).catch(() => {})
  104. await Promise.all([running(agentA), running(agentB)])
  105. expect(agentA.status).toBe('running')
  106. expect(agentB.status).toBe('running')
  107. await harness.ctx.fiber.dispose()
  108. // BOTH agents drained (not still running) — teardown reached quiescence
  109. // across all sessions, not just one.
  110. expect(agentA.status).not.toBe('running')
  111. expect(agentB.status).not.toBe('running')
  112. })
  113. })