Răsfoiți Sursa

test(web): verify subagent.interrupt over the real composition

A browserless keyless-replay e2e boots the shipped web composition, holds
a continuable child's turn open with a replay hang entry, queues a
follow-up and interrupts over plain HTTP, and proves from the real
session state that the turn aborted, the follow-up parked without a new
turn, and a waking send resumed the preserved FIFO order.

Refs #1535
Hypatia May 1 lună în urmă
părinte
comite
82019134b1
3 a modificat fișierele cu 178 adăugiri și 0 ștergeri
  1. 176 0
      apps/web/tests/subagent-interrupt.e2e.ts
  2. 1 0
      apps/web/tsconfig.json
  3. 1 0
      tsconfig.host.json

+ 176 - 0
apps/web/tests/subagent-interrupt.e2e.ts

@@ -0,0 +1,176 @@
+// Web e2e scenario (browserless): the subagent.interrupt RPC against the real
+// composition. A live continuable child holds its model turn open through a
+// replay hang entry; plain HTTP queues a follow-up, interrupts the turn, and
+// proves from the real session state that the turn aborted, the follow-up
+// parked without auto-starting a new turn, and a later waking send resumed the
+// preserved FIFO order. No browser: the RPC surface is the product surface
+// under test, and PR-stacked UI coverage owns the composer interaction.
+import { existsSync } from 'node:fs'
+import { mkdtemp, rm, writeFile } from 'node:fs/promises'
+import { tmpdir } from 'node:os'
+import { join } from 'node:path'
+import { afterAll, beforeAll, describe, expect, it } from 'vitest'
+import { SessionId as sessionId, type SessionId } from '@deepseek-ai/dsh-session'
+import type {} from '@deepseek-ai/dsh-agent'
+import { launchWebScaffold, webSnapshotMode, type WebScaffold } from './scaffold.ts'
+
+const MODE = webSnapshotMode()
+const INITIAL = 'Explain event sourcing in one sentence.'
+const FOLLOWUP = 'Now give the same explanation to a human reader.'
+const WAKING = 'And add one concrete example.'
+
+type RpcResult<T> = { ok: true; value: T } | { ok: false; error: { code: string; message: string } }
+
+/** POST one unary RPC through the real HTTP carrier and unwrap its result. */
+async function rpc<T>(baseUrl: string, method: string, payload: unknown): Promise<RpcResult<T>> {
+  const response = await fetch(`${baseUrl}/api/${method}`, {
+    method: 'POST',
+    headers: { 'content-type': 'application/json' },
+    body: JSON.stringify({
+      type: 'client-request',
+      rpcId: `interrupt-e2e-${method}-${crypto.randomUUID()}`,
+      method,
+      payload,
+    }),
+  })
+  if (!response.ok) throw new Error(`${method} failed over HTTP ${response.status}: ${await response.text()}`)
+  return (await response.json() as { result: RpcResult<T> }).result
+}
+
+/** Poll a synchronous condition (hook-safe; expect.poll is test-body only). */
+async function waitFor(predicate: () => boolean, what: string, timeoutMs = 30_000): Promise<void> {
+  const deadline = Date.now() + timeoutMs
+  while (!predicate()) {
+    if (Date.now() >= deadline) throw new Error(`timed out waiting for ${what}`)
+    await new Promise<void>(resolve => setTimeout(resolve, 10))
+  }
+}
+
+/** One text-only scripted model completion (no tool calls: real tools are mounted). */
+function textCompletion(text: string): object {
+  return {
+    kind: 'chunks',
+    chunks: [
+      { type: 'block-start', index: 0, blockType: 'text' },
+      { type: 'text-delta', index: 0, text },
+      { type: 'block-end', index: 0, block: { type: 'text', text } },
+      { type: 'usage', usage: { inputTokens: 20, outputTokens: 8 } },
+      { type: 'finish', reason: { kind: 'stop' } },
+    ],
+  }
+}
+
+describe.skipIf(MODE === 'record')('web e2e: subagent.interrupt over the real composition', () => {
+  let scaffold: WebScaffold
+  let sidecarRoot: string
+  let readyFile: string
+  let parentId: SessionId
+  let childId: SessionId
+
+  beforeAll(async () => {
+    sidecarRoot = await mkdtemp(join(tmpdir(), 'dsh-web-subagent-interrupt-'))
+    readyFile = join(sidecarRoot, 'hang-ready')
+    // Whole-script replacement: the child's three model calls are the hang
+    // (turn 1, interrupted), the parked follow-up's turn, and the waking turn.
+    // The parent never runs a turn, so the child claims this primary script.
+    await writeFile(join(sidecarRoot, 'replay.override.json'), JSON.stringify([
+      { kind: 'hang', readyFile },
+      textCompletion('resumed response one'),
+      textCompletion('resumed response two'),
+    ]))
+    // Header-only primary fixture: the bare-array override replaces the
+    // derived script entirely; the path only anchors replay installation.
+    await writeFile(
+      join(sidecarRoot, 'session.jsonl'),
+      '{"type":"session","version":0,"id":"primary","createdAt":0}\n',
+    )
+    scaffold = await launchWebScaffold({
+      replayFixture: join(sidecarRoot, 'session.jsonl'),
+      replayOverride: join(sidecarRoot, 'replay.override.json'),
+    })
+
+    // A live parent Agent through the real API; no workspace or browser.
+    const created = await rpc<{ sessionId: string }>(scaffold.baseUrl, 'session.create', {
+      cwd: scaffold.workspaceCwd,
+    })
+    if (!created.ok) throw new Error(`session.create failed: ${created.error.code}`)
+    parentId = sessionId(created.value.sessionId)
+    const parent = scaffold.ctx.agents.get(parentId)
+    if (parent === undefined) throw new Error('created parent session did not publish a live Agent')
+
+    const started = await scaffold.ctx.subagents.startContinuable({
+      provider: 'spawn',
+      label: 'event-sourcing researcher',
+      signal: new AbortController().signal,
+      request: { prompt: [{ type: 'text', text: INITIAL }], parent },
+    })
+    childId = started.childId
+    // The hang entry writes readyFile after its prefix chunks, immediately
+    // before waiting for cancellation: the deterministic "turn is open" gate.
+    await waitFor(() => existsSync(readyFile), 'the held child turn to open')
+  }, 120_000)
+
+  afterAll(async () => {
+    const failures: unknown[] = []
+    await scaffold?.close().catch((error: unknown) => failures.push(error))
+    await rm(sidecarRoot, { recursive: true, force: true }).catch((error: unknown) => failures.push(error))
+    if (failures.length === 1) throw failures[0]
+    if (failures.length > 1) throw new AggregateError(failures, 'subagent interrupt teardown failed')
+  })
+
+  it('parks a queued follow-up on interrupt and resumes it FIFO on a waking send', async () => {
+    // Queue the follow-up while the turn is still open, then interrupt.
+    const queued = await rpc<{ messageId: string }>(scaffold.baseUrl, 'subagent.prompt', {
+      parentSessionId: parentId,
+      childSessionId: childId,
+      mode: 'continuable',
+      content: [{ type: 'text', text: FOLLOWUP }],
+    })
+    expect(queued).toMatchObject({ ok: true })
+
+    const settled = scaffold.whenTurnSettled()
+    const interrupted = await rpc<{ accepted: true }>(scaffold.baseUrl, 'subagent.interrupt', {
+      parentSessionId: parentId,
+      childSessionId: childId,
+      mode: 'continuable',
+    })
+    expect(interrupted).toMatchObject({ ok: true, value: { accepted: true } })
+    // accepted acknowledges the admitted cancel, not quiescence: wait for the
+    // aborted turn/end (the composition's first turn/end) before asserting.
+    expect(await settled).toBe(childId)
+
+    // Parked, not resumed: the Activation stays resident with an idle driver,
+    // the follow-up is retained, and no second turn opened.
+    const child = scaffold.ctx.agents.get(childId)
+    expect(child).toBeDefined()
+    expect(child!.status).toBe('idle')
+    expect(child!.inbox.nextTurn).toHaveLength(1)
+    expect(child!.session.events.filter(event => event.type === 'turn/start')).toHaveLength(1)
+    const lastEnd = child!.session.events.filter(event => event.type === 'turn/end').at(-1)
+    expect((lastEnd)?.data.reason.kind).toBe('aborted')
+
+    // Only an explicit waking send resumes the parked queue, FIFO, then the
+    // child runs both turns to completion and settles.
+    const waking = await rpc<{ messageId: string }>(scaffold.baseUrl, 'subagent.prompt', {
+      parentSessionId: parentId,
+      childSessionId: childId,
+      mode: 'continuable',
+      content: [{ type: 'text', text: WAKING }],
+    })
+    expect(waking).toMatchObject({ ok: true })
+    await expect.poll(() => scaffold.ctx.agents.get(childId), { timeout: 60_000 }).toBeUndefined()
+
+    const loaded = await scaffold.ctx.sessionPersistence.load(childId)
+    // Human-origin messages only: the real composition also injects
+    // runtime-context snapshots as non-user-source messages.
+    const userTexts = loaded.events.flatMap(event => event.type === 'user/message'
+      && event.data.source.kind === 'user'
+      ? event.data.content.flatMap(block => block.type === 'text' ? [block.text] : [])
+      : [])
+    expect(userTexts).toEqual([INITIAL, FOLLOWUP, WAKING])
+    const turnEndKinds = loaded.events
+      .filter(event => event.type === 'turn/end')
+      .map(event => (event).data.reason.kind)
+    expect(turnEndKinds).toEqual(['aborted', 'completed', 'completed'])
+  }, 120_000)
+})

+ 1 - 0
apps/web/tsconfig.json

@@ -67,6 +67,7 @@
     "tests/produced-file-mentions.e2e.ts",
     "tests/goal-bar.e2e.ts",
     "tests/subagent-conversation.e2e.ts",
+    "tests/subagent-interrupt.e2e.ts",
     "tests/sidebar-subagent-activity.e2e.ts",
     "tests/bash-abort-row.e2e.ts",
     "tests/skill-tool-row.e2e.ts",

+ 1 - 0
tsconfig.host.json

@@ -54,6 +54,7 @@
     "apps/web/tests/produced-files.e2e.ts",
     "apps/web/tests/produced-file-mentions.e2e.ts",
     "apps/web/tests/subagent-conversation.e2e.ts",
+    "apps/web/tests/subagent-interrupt.e2e.ts",
     "apps/web/tests/sidebar-subagent-activity.e2e.ts",
     "apps/web/tests/bash-abort-row.e2e.ts",
     "apps/web/tests/skill-tool-row.e2e.ts",