From ab441389b2817fc403fb706d3c1824156aac0463 Mon Sep 17 00:00:00 2001 From: _Kerman Date: Tue, 4 Aug 2026 14:10:07 +0800 Subject: [PATCH] fix(session): remove duplicated turn-end step --- .../client/connection/src/client/fixture.ts | 4 +-- .../src/client/sessions/request-inspection.ts | 17 +++++++---- .../runtime/src/client/sessions/session.ts | 11 ++++++- .../cordis/tool-cordis/src/api-catalog.ts | 2 +- packages/core/agent-loop/src/agent.ts | 2 +- packages/core/session/src/invariant.ts | 4 --- packages/core/session/src/repair.ts | 6 +--- packages/core/session/src/types.ts | 6 ++-- packages/goal/goal-session/src/index.ts | 11 +++---- .../session-persistence/src/coordinator.ts | 29 +++++++------------ 10 files changed, 46 insertions(+), 46 deletions(-) diff --git a/packages/client/connection/src/client/fixture.ts b/packages/client/connection/src/client/fixture.ts index 26af82b2aa..7bbabb6750 100644 --- a/packages/client/connection/src/client/fixture.ts +++ b/packages/client/connection/src/client/fixture.ts @@ -1516,7 +1516,7 @@ export function createFixtureApi(options: FixtureOptions = {}): ApiProxy { }, }) append(sessionId, { type: 'step/end', data: { turn: scenario.turn, step: 1 } }) - append(sessionId, { type: 'turn/end', data: { turn: scenario.turn, step: 1, reason: { kind: 'aborted', reason: { kind: 'user' } }, + append(sessionId, { type: 'turn/end', data: { turn: scenario.turn, reason: { kind: 'aborted', reason: { kind: 'user' } }, } }) retryScenarios.delete(sessionId) setRunning(sessionId, false) @@ -1542,7 +1542,7 @@ export function createFixtureApi(options: FixtureOptions = {}): ApiProxy { }, }) append(sessionId, { type: 'step/end', data: { turn: scenario.turn, step: 1 } }) - append(sessionId, { type: 'turn/end', data: { turn: scenario.turn, step: 1, reason: { kind: 'completed' } } }) + append(sessionId, { type: 'turn/end', data: { turn: scenario.turn, reason: { kind: 'completed' } } }) setRunning(sessionId, false) }, /** Log append WITHOUT the mux emit: a frame lost in transit — history still serves it, the client must repull. */ diff --git a/packages/client/runtime/src/client/sessions/request-inspection.ts b/packages/client/runtime/src/client/sessions/request-inspection.ts index 1962dd970a..6bf22a3132 100644 --- a/packages/client/runtime/src/client/sessions/request-inspection.ts +++ b/packages/client/runtime/src/client/sessions/request-inspection.ts @@ -240,6 +240,7 @@ function promptChange( function deriveRequests(events: readonly SessionEvent[]): readonly RequestView[] { const requests: RequestView[] = [] const ordinaryByStep = new Map() + const lastStepByTurn = new Map() let activeStep: string | undefined let activePrompt: ConversationPromptSnapshot | undefined let activeCompaction: number | undefined @@ -266,6 +267,7 @@ function deriveRequests(events: readonly SessionEvent[]): readonly RequestView[] const { turn, step } = sourceEvent.data const key = requestKey(turn, step) ordinaryByStep.set(key, requests.length) + lastStepByTurn.set(turn, key) requests.push({ purpose: 'assistant', startSeq: sourceEvent.seq, @@ -358,12 +360,15 @@ function deriveRequests(events: readonly SessionEvent[]): readonly RequestView[] }) continue } - if (sourceEvent.type === 'turn/end' && sourceEvent.data.reason.kind === 'error') { - const reason = sourceEvent.data.reason - updateAssistant(ordinaryByStep.get(requestKey(sourceEvent.data.turn, sourceEvent.data.step)), { - status: 'error', - error: displayFailureMessage(reason.error), - }) + if (sourceEvent.type === 'turn/end') { + const lastStep = lastStepByTurn.get(sourceEvent.data.turn) + if (sourceEvent.data.reason.kind === 'error') { + updateAssistant(lastStep === undefined ? undefined : ordinaryByStep.get(lastStep), { + status: 'error', + error: displayFailureMessage(sourceEvent.data.reason.error), + }) + } + lastStepByTurn.delete(sourceEvent.data.turn) continue } diff --git a/packages/client/runtime/src/client/sessions/session.ts b/packages/client/runtime/src/client/sessions/session.ts index 5b934a2deb..ab66a90976 100644 --- a/packages/client/runtime/src/client/sessions/session.ts +++ b/packages/client/runtime/src/client/sessions/session.ts @@ -98,6 +98,8 @@ export class Session implements SessionFace { private readonly transcript = new TranscriptAdapter() private partial: PartialAccumulator | null = null private openCalls = new Map() + /** Last entered step per turn, folded from step/start for terminal error placement. */ + private lastStepByTurn = new Map() /** Operational notices and interrupted-turn terminal nodes merged into the flow by seq. * Derived from window events and rebuilt with partial/openCalls; the transcript is * seq-monotonic, so a plain seq merge preserves event order. */ @@ -796,6 +798,10 @@ export class Session implements SessionFace { } switch (event.type) { case 'turn/start': + this.lastStepByTurn.set(event.data.turn, 0) + return + case 'step/start': + this.lastStepByTurn.set(event.data.turn, event.data.step) return case 'assistant/chunk': { const { turn, step, chunk } = event.data @@ -826,6 +832,7 @@ export class Session implements SessionFace { return } case 'turn/end': { + const lastStep = this.lastStepByTurn.get(event.data.turn) ?? 0 this.turnEnds.set(event.data.turn, event.seq) this.turnEndsRev++ if (event.data.reason.kind === 'aborted') { @@ -841,7 +848,7 @@ export class Session implements SessionFace { seq: event.seq, time: event.time, turn: event.data.turn, - step: event.data.step, + step: lastStep, message: displayFailureMessage(failure), code: failure.code, }) @@ -882,6 +889,7 @@ export class Session implements SessionFace { }) this.derivedRev++ } + this.lastStepByTurn.delete(event.data.turn) return } default: @@ -916,6 +924,7 @@ export class Session implements SessionFace { private rebuildDerivedFromWindow(): void { this.partial = null this.openCalls.clear() + this.lastStepByTurn.clear() this.callsRev++ this.derivedNodes = [] this.derivedRev++ diff --git a/packages/cordis/tool-cordis/src/api-catalog.ts b/packages/cordis/tool-cordis/src/api-catalog.ts index 10cdf09e9a..b676e2607a 100644 --- a/packages/cordis/tool-cordis/src/api-catalog.ts +++ b/packages/cordis/tool-cordis/src/api-catalog.ts @@ -2385,7 +2385,7 @@ export const TYPE_API: readonly TypeApiEntry[] = [ }, { name: 'SessionEventMap', - declaration: 'export interface SessionEventMap {\n \'turn/start\': {\n turn: number;\n };\n \'turn/end\': {\n turn: number;\n step: number;\n reason: TurnEndReason;\n };\n \'step/start\': {\n turn: number;\n step: number;\n };\n \'step/end\': {\n turn: number;\n step: number;\n };\n \'user/message\': UserMessage;\n \'assistant/chunk\': {\n turn: number;\n step: number;\n chunk: StreamChunk;\n };\n \'assistant/message\': {\n turn: number;\n step: number;\n message: AssistantMessage;\n usage?: TokenUsage;\n };\n \'tool/call\': {\n turn: number;\n step: number;\n callId: CallId;\n name: string;\n arguments: string;\n };\n \'tool/result\': {\n turn: number;\n step: number;\n message: ToolResultMessage;\n error?: {\n name: string;\n code: string;\n };\n meta?: JsonValue;\n };\n \'todo/write\': {\n todos: TodoItem[];\n };\n \'request/header\': {\n header: EpochHeader;\n reason: RequestHeaderReason;\n };\n \'request/context\': RequestContext;\n \'session/end-seed\': Record;\n}', + declaration: 'export interface SessionEventMap {\n \'turn/start\': {\n turn: number;\n };\n \'turn/end\': {\n turn: number;\n reason: TurnEndReason;\n };\n \'step/start\': {\n turn: number;\n step: number;\n };\n \'step/end\': {\n turn: number;\n step: number;\n };\n \'user/message\': UserMessage;\n \'assistant/chunk\': {\n turn: number;\n step: number;\n chunk: StreamChunk;\n };\n \'assistant/message\': {\n turn: number;\n step: number;\n message: AssistantMessage;\n usage?: TokenUsage;\n };\n \'tool/call\': {\n turn: number;\n step: number;\n callId: CallId;\n name: string;\n arguments: string;\n };\n \'tool/result\': {\n turn: number;\n step: number;\n message: ToolResultMessage;\n error?: {\n name: string;\n code: string;\n };\n meta?: JsonValue;\n };\n \'todo/write\': {\n todos: TodoItem[];\n };\n \'request/header\': {\n header: EpochHeader;\n reason: RequestHeaderReason;\n };\n \'request/context\': RequestContext;\n \'session/end-seed\': Record;\n}', }, { name: 'SessionEventMetadataFilter', diff --git a/packages/core/agent-loop/src/agent.ts b/packages/core/agent-loop/src/agent.ts index ba18196a0b..2df7e09f2b 100644 --- a/packages/core/agent-loop/src/agent.ts +++ b/packages/core/agent-loop/src/agent.ts @@ -296,7 +296,7 @@ export class ReactLoopAgent implements Agent { } finally { try { // oxlint-disable-next-line typescript/no-non-null-assertion -- every exit assigns a turn ending - this.session.append('turn/end', { turn, step: phase.step, reason: turnEnds! }) + this.session.append('turn/end', { turn, reason: turnEnds! }) } catch (error: unknown) { this.throwError(error) } diff --git a/packages/core/session/src/invariant.ts b/packages/core/session/src/invariant.ts index aedeacc637..f86b43716e 100644 --- a/packages/core/session/src/invariant.ts +++ b/packages/core/session/src/invariant.ts @@ -87,10 +87,6 @@ function validateEvent( if (trace.openStep !== null) { fail(`turn/end ${event.data.turn} while step ${trace.openStep} is still open`) } - const lastStep = trace.nextStep - 1 - if (event.data.step !== lastStep) { - fail(`turn/end ${event.data.turn} expected last step ${lastStep}, got ${event.data.step}`) - } openTurn = null nextTurn += 1 break diff --git a/packages/core/session/src/repair.ts b/packages/core/session/src/repair.ts index f1834e5bf6..1114156c2e 100644 --- a/packages/core/session/src/repair.ts +++ b/packages/core/session/src/repair.ts @@ -46,7 +46,6 @@ export const TOOL_OUTCOME_UNKNOWN = 'TOOL_OUTCOME_UNKNOWN' export function interruptedTurnClosers(events: readonly SessionEvent[]): SessionEvent[] { let openTurn: number | null = null let openStep: number | null = null - let lastStep = 0 // Reset at each turn boundary so earlier calls cannot leak into tail repair. // Assistant blocks register calls; later tool/call events add provenance seqs. const pendingCalls = new Map() @@ -55,18 +54,15 @@ export function interruptedTurnClosers(events: readonly SessionEvent[]): Session case 'turn/start': openTurn = event.data.turn openStep = null - lastStep = 0 pendingCalls.clear() break case 'turn/end': openTurn = null openStep = null - lastStep = 0 pendingCalls.clear() break case 'step/start': openStep = event.data.step - lastStep = event.data.step break case 'step/end': pendingCalls.clear() @@ -151,6 +147,6 @@ export function interruptedTurnClosers(events: readonly SessionEvent[]): Session if (openStep !== null) { closers.push({ type: 'step/end', seq: seq++, time, data: { turn: openTurn, step: openStep } }) } - closers.push({ type: 'turn/end', seq: seq++, time, data: { turn: openTurn, step: lastStep, reason: { kind: 'interrupted' } } }) + closers.push({ type: 'turn/end', seq: seq++, time, data: { turn: openTurn, reason: { kind: 'interrupted' } } }) return closers } diff --git a/packages/core/session/src/types.ts b/packages/core/session/src/types.ts index 0c4c3c822e..4e7026e0a3 100644 --- a/packages/core/session/src/types.ts +++ b/packages/core/session/src/types.ts @@ -196,14 +196,14 @@ export interface SessionEventMap { */ 'turn/start': { turn: number } /** - * Closes turn `turn` after `step`, the last entered step (`0` when none), - * with the {@link TurnEndReason} that ended it. The loop does not await a + * Closes turn `turn` with the {@link TurnEndReason} that ended it. A turn + * with no entered step has no `step/start` or `step/end`. The loop does not await a * flush at turn boundaries: `dsh-session-checkpoint-policy` owns the * per-request durability checkpoint, and consumers that read storage after * `whenIdle()` flush themselves. Success commits the turn; rejection is * reported live and does not prevent later work. */ - 'turn/end': { turn: number; step: number; reason: TurnEndReason } + 'turn/end': { turn: number; reason: TurnEndReason } /** Opens step `step` of turn `turn` — one model call plus the tool executions it requested. */ 'step/start': { turn: number; step: number } /** Closes step `step` of turn `turn`. */ diff --git a/packages/goal/goal-session/src/index.ts b/packages/goal/goal-session/src/index.ts index 01803bcf12..b81d1cb599 100644 --- a/packages/goal/goal-session/src/index.ts +++ b/packages/goal/goal-session/src/index.ts @@ -320,7 +320,9 @@ export function apply(ctx: Context): void { return } if (event.data.reason.kind !== 'aborted') return - if (state.attempt?.phase === 'admitted') state.attempt.cancelled = true + if (state.attempt?.phase === 'claimed' || state.attempt?.phase === 'admitted') { + state.attempt.cancelled = true + } else disarm(state) return default: @@ -372,10 +374,9 @@ export function apply(ctx: Context): void { decision = await next() } catch (error: unknown) { if (signal.aborted) throw error - // A throwing downstream hook drops the whole step proposal: the loop - // returns to idle without a turn, so a still-queued reservation would - // starve every later drive pass. Clear it and let the driver - // reschedule the round. + // A throwing downstream hook drops the whole step proposal. Clear the + // reservation before the balanced no-step turn returns to idle so the + // next drive pass can reschedule the round. state.attempt = undefined requestDrive(state) throw error diff --git a/packages/session-persistence/session-persistence/src/coordinator.ts b/packages/session-persistence/session-persistence/src/coordinator.ts index f1a7daaba6..5ba419f118 100644 --- a/packages/session-persistence/session-persistence/src/coordinator.ts +++ b/packages/session-persistence/session-persistence/src/coordinator.ts @@ -71,7 +71,7 @@ export interface PersistenceBackend { * region strictly below `fromSeq` is limited to seq contiguity — the * service contract scopes this read to the suffix — unless that suffix * contains a supported legacy shape whose normalization needs earlier - * step or message-identity facts, in which case the coordinator falls back + * message-identity facts, in which case the coordinator falls back * to the complete stored prefix. * @param id - persisted session id to resolve. * @param fromSeq - first event seq to include (non-negative safe integer, @@ -214,7 +214,6 @@ function needsLegacyPrefix(event: SessionEvent): boolean { const data = asRecord(event.data) const legacySteeringType: string = 'steering/message' if (event.type === legacySteeringType) return true - if (event.type === 'turn/end' && data !== undefined && !Object.hasOwn(data, 'step')) return true if (data === undefined) return false switch (event.type) { case 'user/message': @@ -270,11 +269,11 @@ function migrateLegacyTurnStartEvent(event: SessionEvent, id: SessionId): Sessio return { ...event, data: { turn: data['turn'] } } as SessionEvent } -/** Upgrade the turn boundary emitted immediately before the loop refactor. */ -function migrateLegacyTurnEndEvent(event: SessionEvent, id: SessionId, lastStep: number): SessionEvent { +/** Upgrade an obsolete turn ending while preserving the latest-master envelope. */ +function migrateLegacyTurnEndEvent(event: SessionEvent, id: SessionId): SessionEvent { if (event.type !== 'turn/end') return event const data = asRecord(event.data) - if (data === undefined || Object.hasOwn(data, 'step')) return event + if (data === undefined) return event const malformed = (): never => { throw new Error(`session "${id}" contains malformed pre-react-loop turn/end at seq ${event.seq}`) } @@ -283,15 +282,16 @@ function migrateLegacyTurnEndEvent(event: SessionEvent, id: SessionId, lastStep: || !hasOnlyKeys(data, ['turn', 'reason']) || reason === undefined || typeof reason['kind'] !== 'string') return malformed() - let currentReason: Record + let currentReason: Record | undefined switch (reason['kind']) { case 'completed': + case 'blocked': case 'max-tokens': case 'interrupted': if (!hasOnlyKeys(reason, ['kind'])) return malformed() - currentReason = { kind: reason['kind'] } - break + return event case 'aborted': + if (Object.hasOwn(reason, 'reason')) return event if (!hasOnlyKeys(reason, ['kind'])) return malformed() currentReason = { kind: 'aborted', reason: { kind: 'legacy' } } break @@ -300,6 +300,7 @@ function migrateLegacyTurnEndEvent(event: SessionEvent, id: SessionId, lastStep: currentReason = { kind: 'aborted', reason: { kind: 'disposed' } } break case 'error': { + if (Object.hasOwn(reason, 'error')) return event if (!Number.isSafeInteger(reason['step']) || (reason['step'] as number) < 0) return malformed() const failure = asRecord(reason['failure']) if (failure !== undefined && hasOnlyKeys(reason, ['kind', 'step', 'failure']) @@ -327,14 +328,13 @@ function migrateLegacyTurnEndEvent(event: SessionEvent, id: SessionId, lastStep: break } default: - return malformed() + return event } return { ...event, data: { ...data, - step: lastStep, reason: currentReason, }, } as SessionEvent @@ -431,16 +431,9 @@ function eventMessageId(event: SessionEvent): PersistedMessageId | undefined { function snapshotStoredEvents(events: readonly SessionEvent[], id: SessionId): SessionEvent[] { assertSupportedEvents(events, id) const messageIds = new Map() - const lastSteps = new Map() return events.map((event) => { - const stepData = event.type === 'step/end' ? asRecord(event.data) : undefined - if (typeof stepData?.['turn'] === 'number' && typeof stepData['step'] === 'number') { - lastSteps.set(stepData['turn'], stepData['step']) - } - const turnData = event.type === 'turn/end' ? asRecord(event.data) : undefined - const lastStep = typeof turnData?.['turn'] === 'number' ? lastSteps.get(turnData['turn']) ?? 0 : 0 const migratedStart = migrateLegacyTurnStartEvent(event, id) - const migratedTurn = migrateLegacyTurnEndEvent(migratedStart, id, lastStep) + const migratedTurn = migrateLegacyTurnEndEvent(migratedStart, id) const migratedSteering = migrateLegacySteeringEvent(migratedTurn, id) const snapshot = snapshotSessionEvent(migrateLegacyMessageEvent(migratedSteering, id, messageIds)) const messageId = eventMessageId(snapshot)