Просмотр исходного кода

fix(agent-loop): stop driver after cancellation

_Kerman 2 месяцев назад
Родитель
Сommit
12a48558f2

+ 9 - 11
packages/core/agent-loop/src/agent.ts

@@ -150,7 +150,7 @@ export class ReactLoopAgent implements Agent {
     try {
       while (await this.turn()) {}
     } catch (_error) {
-      // Admission and turn boundaries report before rethrowing; the driver only contains the rejection.
+      // Reported failures and cancellation are contained at the driver boundary.
     } finally {
       if (this.phase.kind === 'running') {
         this.setPhase({ kind: 'idle', lastTurn: this.phase.turn })
@@ -190,15 +190,14 @@ export class ReactLoopAgent implements Agent {
     const lastTurn = this.phase.kind === 'collecting' ? this.phase.lastTurn : this.phase.turn
     const phase = { kind: 'running' as const, abort, turn: lastTurn, step: 0 }
     this.setPhase(phase)
-    if (signal.aborted) return this.inbox.hasPending
+    signal.throwIfAborted()
     let admission: Admission
     try {
       admission = await this.admit(true)
       if (admission.kind !== 'admitted') return false
       signal.throwIfAborted()
     } catch (error: unknown) {
-      // oxlint-disable-next-line typescript/no-unnecessary-condition -- cancel may abort while admission awaits
-      if (signal.aborted) return this.inbox.hasPending
+      if (signal.aborted) throw error
       this.throwError(error)
     }
     const turn = ++phase.turn
@@ -237,16 +236,15 @@ export class ReactLoopAgent implements Agent {
         if (admission.kind === 'empty' && turnEnds) break
       }
     } catch (error: unknown) {
-      // oxlint-disable-next-line typescript/no-unnecessary-condition -- cancel may abort during any awaited turn operation
       if (signal.aborted) {
         turnEnds = { kind: 'aborted', reason: signal.reason as AgentCancelCause }
-      } else {
-        turnEnds = {
-          kind: 'error',
-          error: error instanceof LlmError ? error.failure : errorChain(error),
-        }
-        this.throwError(error)
+        throw error
       }
+      turnEnds = {
+        kind: 'error',
+        error: error instanceof LlmError ? error.failure : errorChain(error),
+      }
+      this.throwError(error)
     } finally {
       try {
         // oxlint-disable-next-line typescript/no-non-null-assertion -- every exit assigns a turn ending

+ 84 - 28
packages/core/agent-loop/tests/cancel.spec.ts

@@ -73,7 +73,7 @@ describe('Agent.cancel()', () => {
   })
 
   it('cancel({ keepInbox: true }) preserves queued work and emits no discard', async () => {
-    const adapter = new MockAdapter([textResponse('reply')])
+    const adapter = new MockAdapter([textResponse('preserved reply'), textResponse('wake reply')])
     const ctx = await harness(adapter)
     const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
 
@@ -85,11 +85,43 @@ describe('Agent.cancel()', () => {
     agent.cancel({ kind: 'user' }, { keepInbox: true })
     expect(agent.session.events.some(event =>
       event.type === 'agent/inbox/spliced' && event.data.outcome === 'canceled')).toBe(false)
+    await agent.whenIdle()
+    expect(agent.inbox.nextTurn).toHaveLength(1)
+    expect(userTexts(agent)).toEqual([])
+    expect(adapter.requests).toHaveLength(0)
 
     // The preserved item still runs once a later follow-up wakes the driver.
+    const idle = waitForIdle(ctx, agent)
     send(agent, 'wake it')
-    await waitForIdle(ctx, agent)
+    await idle
     expect(userTexts(agent)).toEqual(['preserved', 'wake it'])
+    expect(adapter.requests).toHaveLength(2)
+  })
+
+  it('cancel({ keepInbox: true }) parks queued work after an active turn aborts', async () => {
+    const adapter = new MockAdapter([
+      'hang',
+      textResponse('preserved reply'),
+      textResponse('wake reply'),
+    ])
+    const ctx = await harness(adapter)
+    const agent = ctx.agentLoop.create(SessionId('keep-after-abort'), { provider: 'mock', model: 'mock' })
+
+    send(agent, 'active')
+    await new Promise(resolve => setTimeout(resolve, 30))
+    send(agent, 'preserved')
+    agent.cancel({ kind: 'user' }, { keepInbox: true })
+    await agent.whenIdle()
+
+    expect(userTexts(agent)).toEqual(['active'])
+    expect(agent.inbox.nextTurn).toHaveLength(1)
+    expect(adapter.requests).toHaveLength(1)
+
+    const idle = waitForIdle(ctx, agent)
+    send(agent, 'wake it')
+    await idle
+    expect(userTexts(agent)).toEqual(['active', 'preserved', 'wake it'])
+    expect(adapter.requests).toHaveLength(3)
   })
 
   it('pre-step cancel drops the about-to-start turn (no turn is opened)', async () => {
@@ -195,8 +227,12 @@ describe('Agent.cancel()', () => {
     expect(userTexts(agent)).toEqual(['first', 'later'])
   })
 
-  it('replacement work queued after idle-listener cancellation still runs', async () => {
-    const adapter = new MockAdapter([textResponse('first reply'), textResponse('replacement reply')])
+  it('replacement work queued after idle-listener cancellation waits for another wakeup', async () => {
+    const adapter = new MockAdapter([
+      textResponse('first reply'),
+      textResponse('replacement reply'),
+      textResponse('wake reply'),
+    ])
     const ctx = await harness(adapter)
     const agent = ctx.agentLoop.create(SessionId('idle-listener-post-cancel-send'), { provider: 'mock', model: 'mock' })
 
@@ -216,8 +252,15 @@ describe('Agent.cancel()', () => {
     if (replacementIdle === undefined) throw new Error('idle listener did not register replacement work')
     await replacementIdle
 
-    expect(adapter.requests).toHaveLength(2)
-    expect(userTexts(agent)).toEqual(['first', 'surviving replacement'])
+    expect(adapter.requests).toHaveLength(1)
+    expect(userTexts(agent)).toEqual(['first'])
+    expect(agent.inbox.nextTurn).toHaveLength(1)
+
+    const idle = waitForIdle(ctx, agent)
+    send(agent, 'wake it')
+    await idle
+    expect(adapter.requests).toHaveLength(3)
+    expect(userTexts(agent)).toEqual(['first', 'surviving replacement', 'wake it'])
   })
 
   it('cancel() mid-step aborts the active turn and drops every queued tail item', async () => {
@@ -462,8 +505,7 @@ describe('Agent.cancel()', () => {
     expect(agent.session.events.some(e => e.type === 'turn/start')).toBe(false)
   })
 
-  it('window 2: whenIdle() does NOT resolve early when a running listener cancels then queues replacement work', async () => {
-    // Cancellation must not settle idle while replacement work remains queued.
+  it('a running-listener cancellation parks replacement work until another wakeup', async () => {
     const adapter = new MockAdapter([textResponse('A reply'), textResponse('B reply')])
     const ctx = await harness(adapter)
     const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
@@ -481,16 +523,17 @@ describe('Agent.cancel()', () => {
     await idle
     dispose()
 
-    // whenIdle() resolved only AFTER B's turn ran: B's user message + a turn/end
-    // are in the log, and A was dropped.
-    expect(userTexts(agent)).toContain('B')
-    expect(userTexts(agent)).not.toContain('A')
-    expect(agent.session.events.some(e => e.type === 'turn/end')).toBe(true)
+    expect(userTexts(agent)).toEqual([])
+    expect(agent.inbox.nextTurn).toHaveLength(1)
+
+    const replacementIdle = waitForIdle(ctx, agent)
+    send(agent, 'C')
+    await replacementIdle
+    expect(userTexts(agent)).toEqual(['B', 'C'])
+    expect(agent.session.events.filter(event => event.type === 'turn/end')).toHaveLength(2)
   })
 
-  it('whenIdle() does NOT resolve early when a new prompt is queued during a pre-step cancel', async () => {
-    // The subtle race: a whenIdle() waiter is registered for prompt A; cancel() clears A;
-    // prompt B is queued before the loop resumes from the idle wait.
+  it('a prompt queued during pre-step cancellation waits for another wakeup', async () => {
     const adapter = new MockAdapter([textResponse('A reply'), textResponse('B reply')])
     const ctx = await harness(adapter)
     const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
@@ -500,13 +543,15 @@ describe('Agent.cancel()', () => {
     agent.cancel({ kind: 'user' })     // arms marker, clears A
     send(agent, 'B')           // B races in before the loop resumes
 
-    // whenIdle() must resolve only after B's turn fully ran — by which point B's user message
-    // and a turn/end are in the log.
     await idle
-    expect(userTexts(agent)).toContain('B')
-    expect(agent.session.events.some(e => e.type === 'turn/end')).toBe(true)
-    // A was dropped (never ran); only B's turn is recorded.
-    expect(userTexts(agent)).not.toContain('A')
+    expect(userTexts(agent)).toEqual([])
+    expect(agent.inbox.nextTurn).toHaveLength(1)
+
+    const replacementIdle = waitForIdle(ctx, agent)
+    send(agent, 'C')
+    await replacementIdle
+    expect(userTexts(agent)).toEqual(['B', 'C'])
+    expect(agent.session.events.filter(event => event.type === 'turn/end')).toHaveLength(2)
   })
 
   it("cancel clears the turn's steering — it is not re-enqueued as a fresh turn", async () => {
@@ -537,8 +582,12 @@ describe('Agent.cancel()', () => {
     expect(flat).not.toContain('steer text')
   })
 
-  it('keeps replacement work queued synchronously by an abort observer', async () => {
-    const adapter = new MockAdapter(['hang', textResponse('replacement reply')])
+  it('parks replacement work queued synchronously by an abort observer', async () => {
+    const adapter = new MockAdapter([
+      'hang',
+      textResponse('replacement reply'),
+      textResponse('wake reply'),
+    ])
     const ctx = await harness(adapter)
     const agent = ctx.agentLoop.create(SessionId('abort-observer-replacement'), { provider: 'mock', model: 'mock' })
 
@@ -547,7 +596,7 @@ describe('Agent.cancel()', () => {
     const signal = adapter.requests[0]?.signal
     if (signal === undefined) throw new Error('model request omitted its turn signal')
     signal.addEventListener('abort', () => { send(agent, 'replacement') }, { once: true })
-    const idle = waitForIdle(ctx, agent)
+    const idle = agent.whenIdle()
     agent.cancel({ kind: 'user' })
     await Promise.race([
       idle,
@@ -563,12 +612,19 @@ describe('Agent.cancel()', () => {
       }),
     ])
 
-    expect(adapter.requests).toHaveLength(2)
-    expect(userTexts(agent)).toEqual(['original', 'replacement'])
+    expect(adapter.requests).toHaveLength(1)
+    expect(userTexts(agent)).toEqual(['original'])
+    expect(agent.inbox.nextTurn).toHaveLength(1)
     const reasons = agent.session.events
       .filter(event => event.type === 'turn/end')
       .map(event => event.type === 'turn/end' ? event.data.reason : undefined)
-    expect(reasons).toEqual([{ kind: 'aborted', reason: { kind: 'user' } }, { kind: 'completed' }])
+    expect(reasons).toEqual([{ kind: 'aborted', reason: { kind: 'user' } }])
+
+    const replacementIdle = waitForIdle(ctx, agent)
+    send(agent, 'wake it')
+    await replacementIdle
+    expect(adapter.requests).toHaveLength(3)
+    expect(userTexts(agent)).toEqual(['original', 'replacement', 'wake it'])
   })
 
   it('keeps the first typed cause for an active turn', async () => {

+ 21 - 10
packages/core/agent-loop/tests/contract-regressions.spec.ts

@@ -128,8 +128,11 @@ describe('assistant replay provenance', () => {
 })
 
 describe('abort during tool execution ends the turn', () => {
-  it('records context finalized after a tool-step abort in the next turn', async () => {
-    const adapter = new MockAdapter([toolCallResponse('c1', 'aborter', {})])
+  it('parks context finalized after a tool-step abort until another wakeup', async () => {
+    const adapter = new MockAdapter([
+      toolCallResponse('c1', 'aborter', {}),
+      textResponse('after wake'),
+    ])
     const ctx = await harness(adapter)
     const agent = ctx.agentLoop.create(SessionId('a-abort-injection'), { provider: 'mock', model: 'mock' })
     ctx.tools.register(defineContentToolFixture({
@@ -153,14 +156,20 @@ describe('abort during tool execution ends the turn', () => {
     send(agent, 'go')
     await waitForIdle(ctx, agent)
 
-    const events = [...agent.session.events]
-    expect(events
+    expect(agent.session.events
       .filter(event => event.type === 'tool/result'
         || (event.type === 'user/message' && event.data.source.kind === 'plugin')
         || event.type === 'step/end' || event.type === 'turn/end')
       .map(event => event.type))
-      .toEqual(['tool/result', 'step/end', 'turn/end', 'user/message', 'step/end', 'turn/end'])
-    expect(events
+      .toEqual(['tool/result', 'step/end', 'turn/end'])
+    expect(agent.inbox.nextStep.map(inboxText))
+      .toEqual(['accepted result context after abort'])
+
+    const idle = waitForIdle(ctx, agent)
+    send(agent, 'wake')
+    await idle
+
+    expect(agent.session.events
       .flatMap(event => event.type === 'user/message' && event.data.source.kind === 'plugin'
         ? [event.data.content]
         : []))
@@ -224,7 +233,7 @@ describe('abort during tool execution ends the turn', () => {
       .toBeUndefined()
   })
 
-  it('records result context finalized after disposal cancellation', async () => {
+  it('parks result context finalized after disposal cancellation without opening another turn', async () => {
     const adapter = new MockAdapter([toolCallResponse('c1', 'waiter', {})])
     const ctx = await harness(adapter)
     const started = Promise.withResolvers<undefined>()
@@ -264,9 +273,11 @@ describe('abort during tool execution ends the turn', () => {
       .flatMap(event => event.type === 'user/message' && event.data.source.kind === 'plugin'
         ? [event.data.content]
         : []))
-      .toEqual([
-        [{ type: 'text', text: 'accepted result context during disposal' }],
-      ])
+      .toEqual([])
+    expect(agent.inbox.nextStep.map(inboxText))
+      .toEqual(['accepted result context during disposal'])
+    expect(agent.session.events.filter(event => event.type === 'turn/start'))
+      .toHaveLength(1)
     expect(agent.session.events.find(event => event.type === 'turn/end')?.data.reason)
       .toEqual({ kind: 'aborted', reason: { kind: 'disposed' } })
   })

+ 19 - 5
packages/core/agent-loop/tests/tool-calls.spec.ts

@@ -518,10 +518,10 @@ describe('tool-call scheduler: abort handling', () => {
     ])
   })
 
-  it('stops replenishing after abort, commits started results, and drains accepted additional contexts', async () => {
+  it('stops replenishing after abort, commits started results, and parks accepted additional contexts', async () => {
     const adapter = new MockAdapter([
       multiCall([1, 2, 3, 4].map(n => ({ id: `c${n}`, name: 'p', args: { id: String(n) } }))),
-      textResponse('should never be requested'),
+      textResponse('after wake'),
     ])
     const ctx = await harness(adapter, 2)
     const gated = gatedParallelTool('p')
@@ -566,9 +566,23 @@ describe('tool-call scheduler: abort handling', () => {
     const settled = events(agent).filter(e => e.type === 'tool/result'
       || (e.type === 'user/message' && e.data.source.kind === 'plugin'))
     expect(settled.map(e => e.type))
-      .toEqual(['tool/result', 'tool/result', 'tool/result', 'tool/result', 'user/message', 'user/message'])
-    expect(settled.filter(e => e.type === 'user/message')
-      .map(e => (e.data.content[0] as { text: string }).text))
+      .toEqual(['tool/result', 'tool/result', 'tool/result', 'tool/result'])
+    expect(agent.inbox.nextStep.map(message => message.content[0]))
+      .toEqual([
+        { type: 'text', text: 'ctx-c1' },
+        { type: 'text', text: 'ctx-c2' },
+      ])
+
+    const idle = waitForIdle(ctx, agent)
+    agent.followup(createUserMessage({ content: [{ type: 'text', text: 'wake' }], source: { kind: 'user' } }))
+    await idle
+
+    expect(events(agent).flatMap(e =>
+      e.type === 'user/message'
+        && e.data.source.kind === 'plugin'
+        && e.data.content[0]?.type === 'text'
+        ? [e.data.content[0].text]
+        : []))
       .toEqual(['ctx-c1', 'ctx-c2'])
   })