|
|
@@ -79,7 +79,8 @@ describe('Agent.cancel()', () => {
|
|
|
|
|
|
// send() queues synchronously (status still idle, loop microtask not yet
|
|
|
// resumed). Cancel in that pre-step window: the queued turn must not run.
|
|
|
- send(agent, 'drop me')
|
|
|
+ send(agent, 'drop me first')
|
|
|
+ send(agent, 'drop me second')
|
|
|
agent.cancel('pre-step')
|
|
|
|
|
|
// Give the loop a chance to wake and process the cancel.
|
|
|
@@ -91,6 +92,35 @@ describe('Agent.cancel()', () => {
|
|
|
expect(agent.status).toBe('idle')
|
|
|
})
|
|
|
|
|
|
+ it('disposal from the running notification drops queued work before turn start', async () => {
|
|
|
+ const adapter = new MockAdapter([textResponse('should not run')])
|
|
|
+ const ctx = await harness(adapter)
|
|
|
+ const handle = await ctx.agents.create({
|
|
|
+ sessionId: SessionId('dispose-running-session'),
|
|
|
+ agentOptions: { provider: 'mock', model: 'mock' },
|
|
|
+ })
|
|
|
+ const agent = handle.agent
|
|
|
+
|
|
|
+ const running = Promise.withResolvers<undefined>()
|
|
|
+ let disposalDone: Promise<void> | undefined
|
|
|
+ ctx.on('agent/status', (subject, status) => {
|
|
|
+ if (subject !== agent || status !== 'running') return
|
|
|
+ disposalDone = handle.dispose()
|
|
|
+ running.resolve(undefined)
|
|
|
+ })
|
|
|
+
|
|
|
+ send(agent, 'drop before claim')
|
|
|
+ await running.promise
|
|
|
+ if (disposalDone === undefined) throw new Error('running listener did not start disposal')
|
|
|
+ await disposalDone
|
|
|
+ await driverDone(agent)
|
|
|
+
|
|
|
+ expect(agent.status).toBe('disposed')
|
|
|
+ expect(agent.session.events.some(event => event.type === 'turn/start')).toBe(false)
|
|
|
+ expect(userTexts(agent)).toEqual([])
|
|
|
+ expect(adapter.requests).toHaveLength(0)
|
|
|
+ })
|
|
|
+
|
|
|
it('a whenIdle() waiter registered BEFORE a pre-step cancel resolves (F1 hang guard)', async () => {
|
|
|
const adapter = new MockAdapter([textResponse('x')])
|
|
|
const ctx = await harness(adapter)
|
|
|
@@ -110,7 +140,162 @@ describe('Agent.cancel()', () => {
|
|
|
expect(agent.status).toBe('idle')
|
|
|
})
|
|
|
|
|
|
- it('cancel() mid-step aborts the in-flight model call; the turn ends aborted', async () => {
|
|
|
+ it('cancel() between consecutive turns restores idle and leaves idle steer usable', async () => {
|
|
|
+ const adapter = new MockAdapter([textResponse('first reply'), textResponse('steer reply')])
|
|
|
+ const ctx = await harness(adapter)
|
|
|
+ const agent = ctx.agentLoop.create(SessionId('between-turn-cancel'), { provider: 'mock', model: 'mock' })
|
|
|
+
|
|
|
+ let rejectFirstFlush = true
|
|
|
+ ctx.on('session/flush', (session) => {
|
|
|
+ if (session !== agent.session || !rejectFirstFlush) return
|
|
|
+ rejectFirstFlush = false
|
|
|
+ throw new Error('first flush failed')
|
|
|
+ })
|
|
|
+
|
|
|
+ const cancelled = Promise.withResolvers<undefined>()
|
|
|
+ ctx.on('agent/error', (subject, _turn, _step, error) => {
|
|
|
+ if (subject !== agent || error.message !== 'first flush failed') return
|
|
|
+ // The first hop runs before runLoop resumes from runTurn; the second lands
|
|
|
+ // before its resolved waitForQueued continuation checks cancellation.
|
|
|
+ queueMicrotask(() => {
|
|
|
+ queueMicrotask(() => {
|
|
|
+ agent.cancel('between turns')
|
|
|
+ cancelled.resolve(undefined)
|
|
|
+ })
|
|
|
+ })
|
|
|
+ })
|
|
|
+
|
|
|
+ const statuses: string[] = []
|
|
|
+ ctx.on('agent/status', (subject, status) => {
|
|
|
+ if (subject === agent) statuses.push(status)
|
|
|
+ })
|
|
|
+
|
|
|
+ send(agent, 'first')
|
|
|
+ send(agent, 'queued tail')
|
|
|
+ await cancelled.promise
|
|
|
+
|
|
|
+ expect(agent.status).toBe('idle')
|
|
|
+ expect(statuses).toEqual(['running', 'idle'])
|
|
|
+ expect(adapter.requests).toHaveLength(1)
|
|
|
+ expect(agent.session.events.filter(event => event.type === 'turn/start')).toHaveLength(1)
|
|
|
+ expect(userTexts(agent)).toEqual(['first'])
|
|
|
+
|
|
|
+ let idleResolved = false
|
|
|
+ void agent.whenIdle().then(() => { idleResolved = true })
|
|
|
+ await Promise.resolve()
|
|
|
+ expect(idleResolved).toBe(true)
|
|
|
+
|
|
|
+ const idle = waitForIdle(ctx, agent)
|
|
|
+ agent.steer([{ type: 'text', text: 'idle steer' }])
|
|
|
+ await idle
|
|
|
+
|
|
|
+ expect(statuses).toEqual(['running', 'idle', 'running', 'idle'])
|
|
|
+ expect(adapter.requests).toHaveLength(2)
|
|
|
+ expect(userTexts(agent)).toEqual(['first', 'idle steer'])
|
|
|
+ })
|
|
|
+
|
|
|
+ it('an idle-listener replacement keeps whenIdle pending until the replacement turn finishes', async () => {
|
|
|
+ const adapter = new MockAdapter([textResponse('first reply'), textResponse('replacement reply')])
|
|
|
+ const ctx = await harness(adapter)
|
|
|
+ const agent = ctx.agentLoop.create(SessionId('between-turn-idle-listener'), { provider: 'mock', model: 'mock' })
|
|
|
+
|
|
|
+ let rejectFirstFlush = true
|
|
|
+ ctx.on('session/flush', (session) => {
|
|
|
+ if (session !== agent.session || !rejectFirstFlush) return
|
|
|
+ rejectFirstFlush = false
|
|
|
+ throw new Error('first flush failed')
|
|
|
+ })
|
|
|
+
|
|
|
+ ctx.on('agent/error', (subject, _turn, _step, error) => {
|
|
|
+ if (subject !== agent || error.message !== 'first flush failed') return
|
|
|
+ queueMicrotask(() => {
|
|
|
+ queueMicrotask(() => { agent.cancel('between turns') })
|
|
|
+ })
|
|
|
+ })
|
|
|
+
|
|
|
+ const replacementRegistered = Promise.withResolvers<undefined>()
|
|
|
+ let replacementObservation: Promise<{ status: string; requests: number; turns: number }> | undefined
|
|
|
+ ctx.on('agent/status', (subject, status) => {
|
|
|
+ if (subject !== agent || status !== 'idle' || replacementObservation !== undefined) return
|
|
|
+ send(agent, 'replacement')
|
|
|
+ replacementObservation = agent.whenIdle().then(() => ({
|
|
|
+ status: agent.status,
|
|
|
+ requests: adapter.requests.length,
|
|
|
+ turns: agent.session.events.filter(event => event.type === 'turn/start').length,
|
|
|
+ }))
|
|
|
+ replacementRegistered.resolve(undefined)
|
|
|
+ })
|
|
|
+
|
|
|
+ send(agent, 'first')
|
|
|
+ send(agent, 'cancelled tail')
|
|
|
+ await replacementRegistered.promise
|
|
|
+ if (replacementObservation === undefined) throw new Error('idle listener did not register replacement work')
|
|
|
+
|
|
|
+ await expect(replacementObservation).resolves.toEqual({ status: 'idle', requests: 2, turns: 2 })
|
|
|
+ expect(userTexts(agent)).toEqual(['first', 'replacement'])
|
|
|
+ })
|
|
|
+
|
|
|
+ it('idle-listener cancellation settles its waiter without cancelling later work', async () => {
|
|
|
+ const adapter = new MockAdapter([textResponse('first reply'), textResponse('later reply')])
|
|
|
+ const ctx = await harness(adapter)
|
|
|
+ const agent = ctx.agentLoop.create(SessionId('idle-listener-cancel'), { provider: 'mock', model: 'mock' })
|
|
|
+
|
|
|
+ const replacementRegistered = Promise.withResolvers<undefined>()
|
|
|
+ let replacementObservation: Promise<{ status: string; requests: number; turns: number }> | undefined
|
|
|
+ ctx.on('agent/status', (subject, status) => {
|
|
|
+ if (subject !== agent || status !== 'idle' || replacementObservation !== undefined) return
|
|
|
+ send(agent, 'cancelled replacement')
|
|
|
+ replacementObservation = agent.whenIdle().then(() => ({
|
|
|
+ status: agent.status,
|
|
|
+ requests: adapter.requests.length,
|
|
|
+ turns: agent.session.events.filter(event => event.type === 'turn/start').length,
|
|
|
+ }))
|
|
|
+ agent.cancel('idle listener')
|
|
|
+ replacementRegistered.resolve(undefined)
|
|
|
+ })
|
|
|
+
|
|
|
+ send(agent, 'first')
|
|
|
+ await replacementRegistered.promise
|
|
|
+ if (replacementObservation === undefined) throw new Error('idle listener did not register replacement work')
|
|
|
+
|
|
|
+ await expect(Promise.race([
|
|
|
+ replacementObservation,
|
|
|
+ new Promise((_resolve, reject) => setTimeout(() => { reject(new Error('whenIdle hung after idle-listener cancel')) }, 1000)),
|
|
|
+ ])).resolves.toEqual({ status: 'idle', requests: 1, turns: 1 })
|
|
|
+
|
|
|
+ const idle = waitForIdle(ctx, agent)
|
|
|
+ send(agent, 'later')
|
|
|
+ await idle
|
|
|
+ expect(adapter.requests).toHaveLength(2)
|
|
|
+ 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')])
|
|
|
+ const ctx = await harness(adapter)
|
|
|
+ const agent = ctx.agentLoop.create(SessionId('idle-listener-post-cancel-send'), { provider: 'mock', model: 'mock' })
|
|
|
+
|
|
|
+ const replacementRegistered = Promise.withResolvers<undefined>()
|
|
|
+ let replacementIdle: Promise<void> | undefined
|
|
|
+ ctx.on('agent/status', (subject, status) => {
|
|
|
+ if (subject !== agent || status !== 'idle' || replacementIdle !== undefined) return
|
|
|
+ send(agent, 'cancelled replacement')
|
|
|
+ agent.cancel('idle listener')
|
|
|
+ send(agent, 'surviving replacement')
|
|
|
+ replacementIdle = agent.whenIdle()
|
|
|
+ replacementRegistered.resolve(undefined)
|
|
|
+ })
|
|
|
+
|
|
|
+ send(agent, 'first')
|
|
|
+ await replacementRegistered.promise
|
|
|
+ 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'])
|
|
|
+ })
|
|
|
+
|
|
|
+ it('cancel() mid-step aborts the active turn and drops every queued tail item', async () => {
|
|
|
const adapter = new MockAdapter(['hang'])
|
|
|
const ctx = await harness(adapter)
|
|
|
const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
|
|
|
@@ -121,10 +306,14 @@ describe('Agent.cancel()', () => {
|
|
|
send(agent, 'go')
|
|
|
await new Promise(r => setTimeout(r, 30))
|
|
|
expect(agent.status).toBe('running')
|
|
|
+ send(agent, 'queued tail')
|
|
|
agent.cancel('mid-step')
|
|
|
await waitForIdle(ctx, agent)
|
|
|
|
|
|
expect(reasons).toEqual([{ kind: 'aborted', reason: 'mid-step' }])
|
|
|
+ expect(userTexts(agent)).toEqual(['go'])
|
|
|
+ expect(agent.session.events.filter(event => event.type === 'turn/start')).toHaveLength(1)
|
|
|
+ expect(adapter.requests).toHaveLength(1)
|
|
|
})
|
|
|
|
|
|
it('cancel() with no reason defaults to "cancelled" when aborting an in-flight step', async () => {
|