|
|
@@ -969,6 +969,110 @@ describe('P1-5: a started turn (and any open step) is always closed on a boundar
|
|
|
await waitForIdle(ctx, agent)
|
|
|
expect(boundaryCounts(agent).turnEnd).toBe(2)
|
|
|
})
|
|
|
+
|
|
|
+ it('a throwing session/event listener on the error event still closes the turn (finalizer containment)', async () => {
|
|
|
+ // failTurn appends the `error` event; Session.append pushes it BEFORE
|
|
|
+ // notifying session/event listeners, so a throwing listener leaves `error`
|
|
|
+ // in the log but must NOT abort finalization — `reason` is set before the
|
|
|
+ // append and the throw is contained, so closeTurn(false) still runs and
|
|
|
+ // turn/end is appended (the turn is balanced, not left open).
|
|
|
+ // Plain harness (no invariants oracle): the throwing listener is itself a
|
|
|
+ // session/event subscriber. A finish-error drives the boundary-error path.
|
|
|
+ const errorStream: StreamChunk[] = [{ type: 'finish', reason: { kind: 'error', message: 'provider down' } }]
|
|
|
+ const adapter = new MockAdapter([errorStream, textResponse('turn 2 ok')])
|
|
|
+ const ctx = await harness(adapter)
|
|
|
+ const agent = ctx.agentLoop.create('a-errthrow', { model: 'mock' })
|
|
|
+
|
|
|
+ let threw = false
|
|
|
+ ctx.on('session/event', (_s, event) => {
|
|
|
+ if (!threw && event.type === 'error') { threw = true; throw new Error('boom error-event listener') }
|
|
|
+ })
|
|
|
+
|
|
|
+ send(agent, 'go')
|
|
|
+ await waitForIdle(ctx, agent)
|
|
|
+
|
|
|
+ const e = [...agent.session.events]
|
|
|
+ // The error event is in the log (pushed before the listener threw)…
|
|
|
+ expect(e.some(x => x.type === 'error')).toBe(true)
|
|
|
+ // …and the turn was still closed with the error reason (finalization did not
|
|
|
+ // abort): the last event is turn/end carrying the error reason.
|
|
|
+ const last = e.at(-1)
|
|
|
+ expect(last?.type).toBe('turn/end')
|
|
|
+ expect(last?.type === 'turn/end' && last.data.reason).toMatchObject({ kind: 'error', message: 'provider down' })
|
|
|
+
|
|
|
+ // loop survives: a second turn runs normally.
|
|
|
+ send(agent, 'again')
|
|
|
+ await waitForIdle(ctx, agent)
|
|
|
+ expect(adapter.requests).toHaveLength(2)
|
|
|
+ })
|
|
|
+
|
|
|
+ it('a throwing session/event listener on step/end during finalization still appends turn/end', async () => {
|
|
|
+ // A throwing agent/step-start listener drives the outer catch, which calls
|
|
|
+ // closeStep() during finalization. closeStep appends step/end; a
|
|
|
+ // session/event listener throwing on THAT must not abort the catch before
|
|
|
+ // closeTurn(false) — step/end is already logged (balance holds) and the
|
|
|
+ // throw is contained + surfaced via failTurn, so turn/end is still appended.
|
|
|
+ const adapter = new MockAdapter([textResponse('never reached')])
|
|
|
+ const ctx = await harness(adapter)
|
|
|
+ const agent = ctx.agentLoop.create('a-stependthrow', { model: 'mock' })
|
|
|
+
|
|
|
+ // Open a step, then make the agent/step-start emit throw (boundary throw →
|
|
|
+ // outer catch → closeStep during finalization).
|
|
|
+ ctx.on('agent/step-start', () => { throw new Error('boom step-start') })
|
|
|
+ let threw = false
|
|
|
+ ctx.on('session/event', (_s, event) => {
|
|
|
+ if (!threw && event.type === 'step/end') { threw = true; throw new Error('boom step/end listener') }
|
|
|
+ })
|
|
|
+ const errors: Error[] = []
|
|
|
+ ctx.on('agent/error', (_a, _t, _s, error) => void errors.push(error))
|
|
|
+
|
|
|
+ send(agent, 'go')
|
|
|
+ await waitForIdle(ctx, agent)
|
|
|
+
|
|
|
+ const e = [...agent.session.events]
|
|
|
+ // Both step/end and turn/end are present — finalization ran to completion.
|
|
|
+ expect(e.some(x => x.type === 'step/end')).toBe(true)
|
|
|
+ expect(e.some(x => x.type === 'turn/end')).toBe(true)
|
|
|
+ expect(e.at(-1)?.type).toBe('turn/end')
|
|
|
+ expect(errors.length).toBeGreaterThanOrEqual(1) // surfaced via agent/error
|
|
|
+
|
|
|
+ // loop survives.
|
|
|
+ send(agent, 'again')
|
|
|
+ await waitForIdle(ctx, agent)
|
|
|
+ expect(e.filter(x => x.type === 'turn/start').length).toBeGreaterThanOrEqual(1)
|
|
|
+ })
|
|
|
+
|
|
|
+ it('a throwing session/event listener on turn/end is contained (turn still balanced, loop survives)', async () => {
|
|
|
+ // closeTurn appends turn/end; Session.append pushes it BEFORE notifying
|
|
|
+ // session/event listeners, so a throwing listener leaves turn/end in the log
|
|
|
+ // (the turn is balanced) but must not escape — from the normal-path
|
|
|
+ // closeTurn(true) it would otherwise propagate; the append is contained so
|
|
|
+ // the turn/end emit + loop continue. (A throwing agent/turn-end LISTENER is
|
|
|
+ // a separate, already-tested path; here the session/event append notify is
|
|
|
+ // what throws.)
|
|
|
+ const adapter = new MockAdapter([textResponse('turn 1'), textResponse('turn 2')])
|
|
|
+ const ctx = await harness(adapter)
|
|
|
+ const agent = ctx.agentLoop.create('a-turnendappend', { model: 'mock' })
|
|
|
+
|
|
|
+ let threw = false
|
|
|
+ ctx.on('session/event', (_s, event) => {
|
|
|
+ if (!threw && event.type === 'turn/end') { threw = true; throw new Error('boom turn/end listener') }
|
|
|
+ })
|
|
|
+
|
|
|
+ send(agent, 'go')
|
|
|
+ await waitForIdle(ctx, agent)
|
|
|
+ // turn 1 is balanced despite the throwing turn/end listener.
|
|
|
+ const e1 = [...agent.session.events]
|
|
|
+ expect(e1.filter(x => x.type === 'turn/start')).toHaveLength(1)
|
|
|
+ expect(e1.filter(x => x.type === 'turn/end')).toHaveLength(1)
|
|
|
+ expect(e1.at(-1)?.type).toBe('turn/end')
|
|
|
+
|
|
|
+ // loop survives: a second turn runs to completion.
|
|
|
+ send(agent, 'again')
|
|
|
+ await waitForIdle(ctx, agent)
|
|
|
+ expect(adapter.requests).toHaveLength(2)
|
|
|
+ expect([...agent.session.events].filter(x => x.type === 'turn/end')).toHaveLength(2)
|
|
|
+ })
|
|
|
})
|
|
|
|
|
|
describe('P1-7: tool/result is logged under the originating call.id, not result.callId', () => {
|