From 0eb5c12d496dc1f57e926cc18d85764c2942408d Mon Sep 17 00:00:00 2001 From: Hugo Richard Date: Wed, 5 Aug 2026 23:04:00 +0100 Subject: [PATCH 1/2] fix(eve): terminate turns on cancellation and session failure --- .changeset/eve-terminate-cancelled-turns.md | 9 +++ packages/evlog/src/eve/index.ts | 82 ++++++++++++++++++- packages/evlog/test/eve.test.ts | 90 ++++++++++++++++++++- 3 files changed, 179 insertions(+), 2 deletions(-) create mode 100644 .changeset/eve-terminate-cancelled-turns.md diff --git a/.changeset/eve-terminate-cancelled-turns.md b/.changeset/eve-terminate-cancelled-turns.md new file mode 100644 index 00000000..ac1f7691 --- /dev/null +++ b/.changeset/eve-terminate-cancelled-turns.md @@ -0,0 +1,9 @@ +--- +"evlog": patch +--- + +Emit a wide event for eve turns that end without `turn.completed` or `turn.failed`. + +A turn cancelled by eve produced no wide event, and its logger, accumulator and session slot stayed in memory for good: LRU eviction skips any session that still has an active turn, so nothing ever reclaimed them. `turn.cancelled` now closes the turn on its own terminal path — status `499`, `eve.phase: 'cancelled'`, level `info`, because cancellation is not a failure in eve's model. + +`session.failed` and `session.completed` flush any turn still open for that session and drop its carried-over context, which also ends the indefinite retention of snapshots for finished sessions. diff --git a/packages/evlog/src/eve/index.ts b/packages/evlog/src/eve/index.ts index 6e7b389e..5d4042f4 100644 --- a/packages/evlog/src/eve/index.ts +++ b/packages/evlog/src/eve/index.ts @@ -14,6 +14,9 @@ import { const DEFAULT_MAX_SESSIONS = 256 +/** Client-closed-request status used for turns eve cancelled before a terminal outcome. */ +const CANCELLED_STATUS = 499 + /** Options for {@link defineEvlogHook}. */ export interface EvlogEveOptions extends BaseEvlogOptions { /** Passed to {@link initLogger} on the first hook invocation. */ @@ -305,6 +308,7 @@ export function useLogger(ctx?: EveTurnSessionContext): AuditableLogger { interface EveGlobalState { turnStates: Map activeTurnBySession: Map + sessionTurnIds: Map> sessionSnapshots: Map> sessionPendingActions: Map> sessionApprovals: Map @@ -323,6 +327,7 @@ function getEveGlobalState(): EveGlobalState { host[EVE_GLOBAL_STATE] = { turnStates: new Map(), activeTurnBySession: new Map(), + sessionTurnIds: new Map(), sessionSnapshots: new Map(), sessionPendingActions: new Map(), sessionApprovals: new Map(), @@ -342,6 +347,15 @@ function activeTurnBySession(): Map { return getEveGlobalState().activeTurnBySession } +function sessionTurnIds(): Map> { + return getEveGlobalState().sessionTurnIds +} + +/** Turn ids still open for a session, as a snapshot safe to iterate while finishing. */ +function openTurnIds(sessionId: string): string[] { + return [...(sessionTurnIds().get(sessionId) ?? [])] +} + function sessionSnapshots(): Map> { return getEveGlobalState().sessionSnapshots } @@ -371,6 +385,7 @@ function clearSessionState(sessionId: string): void { sessionRollups().delete(sessionId) sessionPendingActions().delete(sessionId) sessionApprovals().delete(sessionId) + sessionTurnIds().delete(sessionId) } function evictStaleSessions(): void { @@ -409,10 +424,12 @@ function derivePhase( accumulator: TurnAccumulator, httpStatus: number, ): string | undefined { + const eve = ctx.eve as { cancelled?: boolean } | undefined + if (eve?.cancelled) return 'cancelled' const approval = ctx.approval as { status?: string } | undefined if (approval?.status === 'rejected') return 'rejected' if (approval?.status === 'pending' || accumulator.pausedForInput) return 'awaiting-approval' - if (httpStatus >= 400) return 'failed' + if (httpStatus >= 400 && httpStatus !== CANCELLED_STATUS) return 'failed' return undefined } @@ -564,6 +581,9 @@ function getOrCreateTurnState( turnStates().set(key, state) activeTurnBySession().set(sessionId, turnId) + const open = sessionTurnIds().get(sessionId) ?? new Set() + open.add(turnId) + sessionTurnIds().set(sessionId, open) return state } @@ -601,10 +621,28 @@ async function finishTurn( if (activeTurnBySession().get(sessionId) === turnId) { activeTurnBySession().delete(sessionId) } + const open = sessionTurnIds().get(sessionId) + open?.delete(turnId) + if (open?.size === 0) sessionTurnIds().delete(sessionId) pruneEmptySessionMaps(sessionId) } } +/** + * Emit every turn still open for a session. eve ends a session with + * `session.completed` / `session.failed`; a turn left open at that point never + * received its own terminal event, so without this it would neither be emitted + * nor released. + */ +async function finishOpenTurns( + sessionId: string, + opts: { status?: number; error?: Error }, +): Promise { + for (const turnId of openTurnIds(sessionId)) { + await finishTurn(sessionId, turnId, opts) + } +} + function runSafe(fn: () => void | Promise): void { void (async () => { try { @@ -821,6 +859,47 @@ export function defineEvlogHook(options: EvlogEveOptions = {}): HookDefinition { } }, + async 'turn.cancelled'(event, ctx) { + try { + const state = getTurnState(ctx.session.id, event.data.turnId) + state?.logger.set({ eve: { cancelled: true } }) + await finishTurn(ctx.session.id, event.data.turnId, { status: CANCELLED_STATUS }) + } catch (err) { + console.error('[evlog] eve hook handler failed:', err) + } + }, + + async 'session.completed'(_event, ctx) { + try { + await finishOpenTurns(ctx.session.id, { status: 200 }) + clearSessionState(ctx.session.id) + } catch (err) { + console.error('[evlog] eve hook handler failed:', err) + } + }, + + async 'session.failed'(event, ctx) { + try { + const error = new Error(event.data.message) + error.name = event.data.code + for (const turnId of openTurnIds(ctx.session.id)) { + getTurnState(ctx.session.id, turnId)?.logger.set({ + eve: { + failure: { + code: event.data.code, + message: event.data.message, + ...(event.data.details ? { details: event.data.details } : {}), + }, + }, + }) + await finishTurn(ctx.session.id, turnId, { error, status: 500 }) + } + clearSessionState(ctx.session.id) + } catch (err) { + console.error('[evlog] eve hook handler failed:', err) + } + }, + async 'turn.failed'(event, ctx) { try { const state = getTurnState(ctx.session.id, event.data.turnId) @@ -855,6 +934,7 @@ export function resetEvlogEveForTests(): void { clearAsyncLocalStorage(turnLoggerStorage) turnStates().clear() activeTurnBySession().clear() + sessionTurnIds().clear() sessionSnapshots().clear() sessionPendingActions().clear() sessionApprovals().clear() diff --git a/packages/evlog/test/eve.test.ts b/packages/evlog/test/eve.test.ts index 6148e2ae..10c0610a 100644 --- a/packages/evlog/test/eve.test.ts +++ b/packages/evlog/test/eve.test.ts @@ -43,6 +43,7 @@ async function runTurn( options: { turnId?: string fail?: boolean + cancel?: boolean steps?: number toolResults?: Array<{ toolName: string @@ -181,7 +182,12 @@ async function runTurn( }, ctx) } - if (options.fail) { + if (options.cancel) { + await events['turn.cancelled']!({ + type: 'turn.cancelled', + data: { sequence: 99, turnId }, + }, ctx) + } else if (options.fail) { await events['turn.failed']!({ type: 'turn.failed', data: { @@ -301,6 +307,88 @@ describe('evlog/eve', () => { }) }) + it('emits a cancelled turn as a non-error wide event', async () => { + const spies = createPipelineSpies() + const hook = defineEvlogHook({ drain: spies.drain }) + + await runTurn(hook, { cancel: true }) + + await waitForDrainCalls(spies.drain) + const event = findEventViaDrain(spies.drain, () => true) + expect(event?.status).toBe(499) + expect(event?.level).toBe('info') + expect(event?.eve).toMatchObject({ phase: 'cancelled', cancelled: true }) + }) + + it('releases turn state after a cancelled turn', async () => { + const spies = createPipelineSpies() + const hook = defineEvlogHook({ drain: spies.drain }) + + await runTurn(hook, { cancel: true }) + + expect(() => useLogger(toolContext())).toThrow(/could not find a logger/) + + await runTurn(hook, { turnId: TURN_ID_1 }) + await waitForDrainCalls(spies.drain, 2) + expect(findEventViaDrain(spies.drain, e => e.path?.includes(TURN_ID_1))).toBeDefined() + }) + + it('flushes an in-flight turn when the session fails without turn.failed', async () => { + const spies = createPipelineSpies() + const hook = defineEvlogHook({ drain: spies.drain }) + const ctx = hookContext() + + hook.events!['turn.started']!({ + type: 'turn.started', + data: { sequence: 0, turnId: TURN_ID }, + }, ctx) + + await hook.events!['session.failed']!({ + type: 'session.failed', + data: { code: 'SESSION_ERROR', message: 'session exploded', sessionId: SESSION_ID }, + }, ctx) + + await waitForDrainCalls(spies.drain) + const event = findEventViaDrain(spies.drain, () => true) + expect(event?.status).toBe(500) + expect(event?.level).toBe('error') + expect(event?.eve).toMatchObject({ + failure: { code: 'SESSION_ERROR', message: 'session exploded' }, + }) + expect(() => useLogger(toolContext())).toThrow(/could not find a logger/) + }) + + it('drops session context once the session completes', async () => { + const spies = createPipelineSpies() + const hook = defineEvlogHook({ drain: spies.drain }) + const ctx = hookContext() + + hook.events!['turn.started']!({ + type: 'turn.started', + data: { sequence: 0, turnId: TURN_ID }, + }, ctx) + useLogger(toolContext()).set({ customer: { slug: 'acme' } }) + await hook.events!['turn.completed']!({ + type: 'turn.completed', + data: { sequence: 1, turnId: TURN_ID }, + }, ctx) + + await hook.events!['session.completed']!({ type: 'session.completed' } as never, ctx) + + hook.events!['turn.started']!({ + type: 'turn.started', + data: { sequence: 0, turnId: TURN_ID_1 }, + }, ctx) + await hook.events!['turn.completed']!({ + type: 'turn.completed', + data: { sequence: 1, turnId: TURN_ID_1 }, + }, ctx) + + await waitForDrainCalls(spies.drain, 2) + const secondTurn = findEventViaDrain(spies.drain, e => e.path?.includes(TURN_ID_1)) + expect(secondTurn?.customer).toBeUndefined() + }) + it('does not throw when an internal handler fails', async () => { const hook = defineEvlogHook({ enrich: () => { From be732165ade621cb72192852666a86a865401866 Mon Sep 17 00:00:00 2001 From: Hugo Richard Date: Wed, 5 Aug 2026 23:41:52 +0100 Subject: [PATCH 2/2] fix(eve): clear session state even when a turn fails to finish --- packages/evlog/src/eve/index.ts | 25 ++++++++++++----- packages/evlog/test/eve.test.ts | 49 ++++++++++++++++++++++++++++++--- 2 files changed, 63 insertions(+), 11 deletions(-) diff --git a/packages/evlog/src/eve/index.ts b/packages/evlog/src/eve/index.ts index 5d4042f4..44571e3f 100644 --- a/packages/evlog/src/eve/index.ts +++ b/packages/evlog/src/eve/index.ts @@ -633,13 +633,23 @@ async function finishTurn( * `session.completed` / `session.failed`; a turn left open at that point never * received its own terminal event, so without this it would neither be emitted * nor released. + * + * Each turn is finished independently: a user `keep` callback that throws + * rejects that turn's `finish`, and must not take the remaining turns with it. */ async function finishOpenTurns( sessionId: string, opts: { status?: number; error?: Error }, + decorate?: (state: TurnState) => void, ): Promise { for (const turnId of openTurnIds(sessionId)) { - await finishTurn(sessionId, turnId, opts) + try { + const state = getTurnState(sessionId, turnId) + if (state && decorate) decorate(state) + await finishTurn(sessionId, turnId, opts) + } catch (err) { + console.error('[evlog] eve hook handler failed:', err) + } } } @@ -872,9 +882,10 @@ export function defineEvlogHook(options: EvlogEveOptions = {}): HookDefinition { async 'session.completed'(_event, ctx) { try { await finishOpenTurns(ctx.session.id, { status: 200 }) - clearSessionState(ctx.session.id) } catch (err) { console.error('[evlog] eve hook handler failed:', err) + } finally { + clearSessionState(ctx.session.id) } }, @@ -882,8 +893,8 @@ export function defineEvlogHook(options: EvlogEveOptions = {}): HookDefinition { try { const error = new Error(event.data.message) error.name = event.data.code - for (const turnId of openTurnIds(ctx.session.id)) { - getTurnState(ctx.session.id, turnId)?.logger.set({ + await finishOpenTurns(ctx.session.id, { error, status: 500 }, (state) => { + state.logger.set({ eve: { failure: { code: event.data.code, @@ -892,11 +903,11 @@ export function defineEvlogHook(options: EvlogEveOptions = {}): HookDefinition { }, }, }) - await finishTurn(ctx.session.id, turnId, { error, status: 500 }) - } - clearSessionState(ctx.session.id) + }) } catch (err) { console.error('[evlog] eve hook handler failed:', err) + } finally { + clearSessionState(ctx.session.id) } }, diff --git a/packages/evlog/test/eve.test.ts b/packages/evlog/test/eve.test.ts index 10c0610a..776c0f3a 100644 --- a/packages/evlog/test/eve.test.ts +++ b/packages/evlog/test/eve.test.ts @@ -368,13 +368,14 @@ describe('evlog/eve', () => { data: { sequence: 0, turnId: TURN_ID }, }, ctx) useLogger(toolContext()).set({ customer: { slug: 'acme' } }) - await hook.events!['turn.completed']!({ - type: 'turn.completed', - data: { sequence: 1, turnId: TURN_ID }, - }, ctx) await hook.events!['session.completed']!({ type: 'session.completed' } as never, ctx) + await waitForDrainCalls(spies.drain) + const openTurn = findEventViaDrain(spies.drain, e => e.path?.includes(TURN_ID)) + expect(openTurn?.status).toBe(200) + expect(() => useLogger(toolContext())).toThrow(/could not find a logger/) + hook.events!['turn.started']!({ type: 'turn.started', data: { sequence: 0, turnId: TURN_ID_1 }, @@ -389,6 +390,46 @@ describe('evlog/eve', () => { expect(secondTurn?.customer).toBeUndefined() }) + it('finishes the remaining turns and clears session state when one turn fails to finish', async () => { + const spies = createPipelineSpies() + const hook = defineEvlogHook({ + drain: spies.drain, + keep: (tail) => { + if (tail.path?.endsWith(TURN_ID)) throw new Error('keep exploded') + }, + }) + const ctx = hookContext() + + hook.events!['turn.started']!({ + type: 'turn.started', + data: { sequence: 0, turnId: TURN_ID }, + }, ctx) + useLogger(toolContext()).set({ customer: { slug: 'acme' } }) + hook.events!['turn.started']!({ + type: 'turn.started', + data: { sequence: 1, turnId: TURN_ID_1 }, + }, ctx) + + await hook.events!['session.completed']!({ type: 'session.completed' } as never, ctx) + + await waitForDrainCalls(spies.drain) + expect(findEventViaDrain(spies.drain, e => e.path?.endsWith(TURN_ID))).toBeUndefined() + expect(findEventViaDrain(spies.drain, e => e.path?.endsWith(TURN_ID_1))).toBeDefined() + + hook.events!['turn.started']!({ + type: 'turn.started', + data: { sequence: 2, turnId: 'turn_2' }, + }, ctx) + await hook.events!['turn.completed']!({ + type: 'turn.completed', + data: { sequence: 3, turnId: 'turn_2' }, + }, ctx) + + await waitForDrainCalls(spies.drain, 2) + const thirdTurn = findEventViaDrain(spies.drain, e => e.path?.endsWith('turn_2')) + expect(thirdTurn?.customer).toBeUndefined() + }) + it('does not throw when an internal handler fails', async () => { const hook = defineEvlogHook({ enrich: () => {