subagent-interrupt-ui.e2e.ts 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271
  1. // Web e2e scenario: the composer's primary action interrupts a running
  2. // continuable child. The child holds its model turn open through a replay
  3. // hang entry; the browser proves the single primary Send/Stop toggle, the
  4. // parent-offline disabled-input-with-Stop composer, the subagent.interrupt
  5. // (never session.cancel) transport, the parked follow-up, and the FIFO resume
  6. // on a waking send.
  7. //
  8. // Replay-binding note: only the PRIMARY script can hang, and scripts bind by
  9. // first-call order, so the child issues the composition's first model call
  10. // (claiming the overridden primary) and the parent's one UI prompt — needed
  11. // so the non-blank parent renders its header catalog — binds to a derived
  12. // child fixture afterwards.
  13. import { existsSync } from 'node:fs'
  14. import { mkdtemp, readFile, rm, writeFile } from 'node:fs/promises'
  15. import { tmpdir } from 'node:os'
  16. import { join } from 'node:path'
  17. import { fileURLToPath } from 'node:url'
  18. import type { Browser, Page } from 'playwright'
  19. import { chromium } from 'playwright'
  20. import { afterAll, beforeAll, describe, expect, it, onTestFailed } from 'vitest'
  21. import type { SessionEvent, SessionId } from '@deepseek-ai/dsh-session'
  22. import type {} from '@deepseek-ai/dsh-agent'
  23. import {
  24. acknowledgeReloadConnectionLoss, assertFixtureInventory, captureStableAria, compareOrRefreshGolden,
  25. launchWebScaffold, watchConsole, webSnapshotMode, type WebScaffold,
  26. } from './scaffold.ts'
  27. import { connectFreshWorkspace, newEnglishPage, saveFailureShot } from './support.ts'
  28. const BASE_FIXTURE = fileURLToPath(new URL('./snapshots/live-interactions/session.jsonl', import.meta.url))
  29. const SNAPSHOT_DIR = fileURLToPath(new URL('./snapshots/subagent-interrupt', import.meta.url))
  30. const OFFLINE_COMPOSER_EXPECTED = join(SNAPSHOT_DIR, 'offline-composer.expected.md')
  31. const MODE = webSnapshotMode()
  32. const LABEL = 'event-sourcing researcher'
  33. const INITIAL = 'Explain event sourcing in one sentence.'
  34. const FOLLOWUP = 'Now give the same explanation to a human reader.'
  35. const WAKING = 'And add one concrete example.'
  36. const PARKED_ANSWER = 'parked follow-up answer'
  37. const WAKING_ANSWER = 'waking answer'
  38. /** Poll a synchronous condition (hook-safe; expect.poll is test-body only). */
  39. async function waitFor(predicate: () => boolean, what: string, timeoutMs = 30_000): Promise<void> {
  40. const deadline = Date.now() + timeoutMs
  41. while (!predicate()) {
  42. if (Date.now() >= deadline) throw new Error(`timed out waiting for ${what}`)
  43. await new Promise<void>(resolve => setTimeout(resolve, 10))
  44. }
  45. }
  46. /** One text-only scripted model completion (no tool calls: real tools are mounted). */
  47. function textCompletion(text: string): object {
  48. return {
  49. kind: 'chunks',
  50. chunks: [
  51. { type: 'block-start', index: 0, blockType: 'text' },
  52. { type: 'text-delta', index: 0, text },
  53. { type: 'block-end', index: 0, block: { type: 'text', text } },
  54. { type: 'usage', usage: { inputTokens: 20, outputTokens: 8 } },
  55. { type: 'finish', reason: { kind: 'stop' } },
  56. ],
  57. }
  58. }
  59. describe.skipIf(MODE === 'record')('web e2e: composer interrupt for a running continuable child', () => {
  60. let scaffold: WebScaffold
  61. let browser: Browser
  62. let page: Page
  63. let sidecarRoot: string
  64. let childId: SessionId
  65. let tripwire: ReturnType<typeof watchConsole>
  66. const apiCalls: string[] = []
  67. beforeAll(async () => {
  68. sidecarRoot = await mkdtemp(join(tmpdir(), 'dsh-web-subagent-interrupt-ui-'))
  69. const readyFile = join(sidecarRoot, 'hang-ready')
  70. // The child claims this whole-script replacement: held turn 1, then the
  71. // parked follow-up and waking turns.
  72. await writeFile(join(sidecarRoot, 'replay.override.json'), JSON.stringify([
  73. { kind: 'hang', readyFile },
  74. textCompletion(PARKED_ANSWER),
  75. textCompletion(WAKING_ANSWER),
  76. ]))
  77. await writeFile(
  78. join(sidecarRoot, 'session.jsonl'),
  79. '{"type":"session","version":0,"id":"primary","createdAt":0}\n',
  80. )
  81. // The parent's one prompted turn replays this recorded single text-only
  82. // call (binding is positional, not lineage-aware).
  83. const parentTurnPath = join(sidecarRoot, 'parent-turn.jsonl')
  84. const base = await readFile(BASE_FIXTURE, 'utf8')
  85. const [header, ...events] = base.trimEnd().split('\n')
  86. if (header === undefined) throw new Error('base replay fixture has no header')
  87. await writeFile(parentTurnPath, [
  88. header
  89. .replace('"id":"{{sessionId}}"', '"id":"recorded-parent-turn"')
  90. .replace(/"createdAt":\d+/, '"createdAt":1784998084442'),
  91. ...events,
  92. '',
  93. ].join('\n'))
  94. scaffold = await launchWebScaffold({
  95. replayFixture: join(sidecarRoot, 'session.jsonl'),
  96. replayOverride: join(sidecarRoot, 'replay.override.json'),
  97. replayChildFixtures: [parentTurnPath],
  98. })
  99. browser = await chromium.launch()
  100. page = await newEnglishPage(browser)
  101. page.on('request', (request) => {
  102. const path = new URL(request.url()).pathname
  103. if (path.startsWith('/api/')) apiCalls.push(path)
  104. })
  105. tripwire = watchConsole(page)
  106. await page.goto(scaffold.baseUrl, { waitUntil: 'load' })
  107. await page.waitForSelector('[class*="frame"]', { timeout: 30_000 })
  108. await connectFreshWorkspace(page, scaffold.workspaceCwd)
  109. const parent = scaffold.ctx.agents.roots()[0]
  110. if (parent === undefined) throw new Error('fresh workspace did not publish its parent Agent')
  111. // The child's first model call claims the primary override and holds.
  112. const started = await scaffold.ctx.subagents.startContinuable({
  113. provider: 'spawn',
  114. label: LABEL,
  115. signal: new AbortController().signal,
  116. request: { prompt: [{ type: 'text', text: INITIAL }], parent },
  117. })
  118. childId = started.childId
  119. await waitFor(() => existsSync(readyFile), 'the held child turn to open')
  120. // One prompted parent turn makes the parent non-blank so the session
  121. // header (and its subagent catalog action) renders.
  122. const parentSettled = scaffold.whenTurnSettled()
  123. const parentInput = page.locator('textarea:enabled').first()
  124. await parentInput.fill('Ask a research subagent to explain event sourcing.')
  125. await parentInput.press('Enter')
  126. expect(await parentSettled).toBe(parent.id)
  127. // Reload onto the restart baseline (the proven route to a freshly
  128. // discovered catalog), with the child still live and running host-side.
  129. const warningStart = tripwire.warnings.length
  130. await page.reload({ waitUntil: 'load' })
  131. await page.waitForSelector('[class*="frame"]', { timeout: 30_000 })
  132. await page.getByRole('button', { name: /1 subagent/ }).waitFor({ timeout: 15_000 })
  133. acknowledgeReloadConnectionLoss(tripwire, warningStart)
  134. expect(scaffold.ctx.agents.get(childId)?.status).toBe('running')
  135. }, 120_000)
  136. afterAll(async () => {
  137. const failures: unknown[] = []
  138. await browser?.close().catch((error: unknown) => failures.push(error))
  139. await scaffold?.close().catch((error: unknown) => failures.push(error))
  140. if (sidecarRoot !== undefined) {
  141. await rm(sidecarRoot, { recursive: true, force: true })
  142. .catch((error: unknown) => failures.push(error))
  143. }
  144. if (failures.length === 1) throw failures[0]
  145. if (failures.length > 1) throw new AggregateError(failures, 'subagent interrupt UI teardown failed')
  146. })
  147. it('locks input but keeps the same primary Stop when the parent is offline', async () => {
  148. onTestFailed(() => saveFailureShot(page, 'web-e2e-subagent-interrupt-offline'))
  149. // Simulate a parent that went offline: the catalog delivers
  150. // parentAvailable: false while the child Activation stays live (the
  151. // interrupt RPC itself needs no live parent — PR 1's host coverage).
  152. const pattern = '**/api/subagent.list'
  153. await page.route(pattern, async (route) => {
  154. const response = await route.fetch()
  155. const body = await response.json() as {
  156. result: { ok: true; value: { parentAvailable: boolean } } | { ok: false }
  157. }
  158. if (body.result.ok) body.result.value.parentAvailable = false
  159. await route.fulfill({ response, json: body })
  160. })
  161. try {
  162. await page.getByRole('button', { name: /1 subagent/ }).click()
  163. await page.getByRole('treeitem', { name: new RegExp(LABEL) }).click()
  164. const input = page.getByRole('textbox', {
  165. name: 'Parent session offline; sending is unavailable but you can still stop the run',
  166. })
  167. await input.waitFor({ timeout: 15_000 })
  168. expect(await input.isDisabled()).toBe(true)
  169. // Still exactly one primary action, and it is an enabled Stop.
  170. const stop = page.getByRole('button', { name: 'Stop generating' })
  171. expect(await stop.count()).toBe(1)
  172. expect(await stop.isEnabled()).toBe(true)
  173. expect(await page.getByRole('button', { name: 'Send message' }).count()).toBe(0)
  174. await compareOrRefreshGolden(
  175. OFFLINE_COMPOSER_EXPECTED,
  176. await captureStableAria(page, '[class*="centerCol"]', scaffold.workspaceCwd),
  177. MODE,
  178. )
  179. } finally {
  180. await page.unroute(pattern)
  181. }
  182. }, 60_000)
  183. it('interrupts through subagent.interrupt, parks the follow-up, and resumes it FIFO', async () => {
  184. onTestFailed(() => saveFailureShot(page, 'web-e2e-subagent-interrupt-flow'))
  185. // Reselect the child with the truthful catalog: parent available again.
  186. await page.getByRole('navigation', { name: 'Session hierarchy' })
  187. .getByRole('button').first().click()
  188. await page.getByRole('button', { name: /1 subagent/ }).click()
  189. await page.getByRole('treeitem', { name: new RegExp(LABEL) }).click()
  190. const input = page.getByRole('textbox', { name: 'Message the agent' })
  191. await input.waitFor({ timeout: 15_000 })
  192. expect(await input.isDisabled()).toBe(false)
  193. // Queue a follow-up while the turn is open; the primary stays Stop.
  194. const promptResponse = page.waitForResponse(response =>
  195. new URL(response.url()).pathname === '/api/subagent.prompt')
  196. await input.fill(FOLLOWUP)
  197. await input.press('Enter')
  198. expect(((await (await promptResponse).json()) as { result: { ok: boolean } }).result)
  199. .toMatchObject({ ok: true })
  200. const aborted = new Promise<void>((resolve, reject) => {
  201. const timer = setTimeout(() => {
  202. off()
  203. reject(new Error('interrupt did not reach an aborted turn/end'))
  204. }, 30_000)
  205. const off = scaffold.ctx.on('session/event', (session: { id: SessionId }, event: SessionEvent) => {
  206. if (session.id !== childId || event.type !== 'turn/end') return
  207. clearTimeout(timer)
  208. off()
  209. if (event.data.reason.kind === 'aborted') resolve()
  210. else reject(new Error(`expected an aborted turn/end, got ${event.data.reason.kind}`))
  211. })
  212. })
  213. const stop = page.getByRole('button', { name: 'Stop generating' })
  214. expect(await stop.count()).toBe(1)
  215. const interruptResponse = page.waitForResponse(response =>
  216. new URL(response.url()).pathname === '/api/subagent.interrupt')
  217. await stop.click()
  218. expect(((await (await interruptResponse).json()) as {
  219. result: { ok: boolean; value?: { accepted: boolean } }
  220. }).result).toMatchObject({ ok: true, value: { accepted: true } })
  221. // The addressed child stops through its own RPC, never the generic one.
  222. expect(apiCalls.filter(path => path === '/api/session.cancel')).toEqual([])
  223. await aborted
  224. // Parked: the Activation stays resident and idle with the retained
  225. // follow-up; the primary returns to Send without a new turn starting.
  226. await expect.poll(() => scaffold.ctx.agents.get(childId)?.status, { timeout: 15_000 }).toBe('idle')
  227. const child = scaffold.ctx.agents.get(childId)
  228. expect(child).toBeDefined()
  229. expect(child!.inbox.nextTurn).toHaveLength(1)
  230. expect(child!.session.events.filter(event => event.type === 'turn/start')).toHaveLength(1)
  231. await page.getByRole('button', { name: 'Send message' }).waitFor({ timeout: 15_000 })
  232. // Only the waking send resumes the parked queue, FIFO, to settlement.
  233. await input.fill(WAKING)
  234. await input.press('Enter')
  235. await expect.poll(() => page.getByText(PARKED_ANSWER, { exact: true }).count(), { timeout: 30_000 }).toBe(1)
  236. await expect.poll(() => page.getByText(WAKING_ANSWER, { exact: true }).count(), { timeout: 30_000 }).toBe(1)
  237. await expect.poll(() => scaffold.ctx.agents.get(childId), { timeout: 60_000 }).toBeUndefined()
  238. const loaded = await scaffold.ctx.sessionPersistence.load(childId)
  239. const userTexts = loaded.events.flatMap(event => event.type === 'user/message'
  240. && event.data.source.kind === 'user'
  241. ? event.data.content.flatMap(block => block.type === 'text' ? [block.text] : [])
  242. : [])
  243. expect(userTexts).toEqual([INITIAL, FOLLOWUP, WAKING])
  244. const turnEndKinds = loaded.events
  245. .filter(event => event.type === 'turn/end')
  246. .map(event => event.data.reason.kind)
  247. expect(turnEndKinds).toEqual(['aborted', 'completed', 'completed'])
  248. expect(tripwire.pageErrors).toEqual([])
  249. }, 120_000)
  250. it('keeps its snapshot inventory closed', async () => {
  251. await assertFixtureInventory(SNAPSHOT_DIR, ['offline-composer.expected.md'])
  252. })
  253. })