fix(session): remove duplicated turn-end step

This commit is contained in:
_Kerman
2026-08-04 14:10:07 +08:00
parent e874910a76
commit ab441389b2
10 changed files with 46 additions and 46 deletions
@@ -1516,7 +1516,7 @@ export function createFixtureApi(options: FixtureOptions = {}): ApiProxy {
}, },
}) })
append(sessionId, { type: 'step/end', data: { turn: scenario.turn, step: 1 } }) 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) retryScenarios.delete(sessionId)
setRunning(sessionId, false) 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: '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) setRunning(sessionId, false)
}, },
/** Log append WITHOUT the mux emit: a frame lost in transit — history still serves it, the client must repull. */ /** Log append WITHOUT the mux emit: a frame lost in transit — history still serves it, the client must repull. */
@@ -240,6 +240,7 @@ function promptChange(
function deriveRequests(events: readonly SessionEvent[]): readonly RequestView[] { function deriveRequests(events: readonly SessionEvent[]): readonly RequestView[] {
const requests: RequestView[] = [] const requests: RequestView[] = []
const ordinaryByStep = new Map<string, number>() const ordinaryByStep = new Map<string, number>()
const lastStepByTurn = new Map<number, string>()
let activeStep: string | undefined let activeStep: string | undefined
let activePrompt: ConversationPromptSnapshot | undefined let activePrompt: ConversationPromptSnapshot | undefined
let activeCompaction: number | undefined let activeCompaction: number | undefined
@@ -266,6 +267,7 @@ function deriveRequests(events: readonly SessionEvent[]): readonly RequestView[]
const { turn, step } = sourceEvent.data const { turn, step } = sourceEvent.data
const key = requestKey(turn, step) const key = requestKey(turn, step)
ordinaryByStep.set(key, requests.length) ordinaryByStep.set(key, requests.length)
lastStepByTurn.set(turn, key)
requests.push({ requests.push({
purpose: 'assistant', purpose: 'assistant',
startSeq: sourceEvent.seq, startSeq: sourceEvent.seq,
@@ -358,12 +360,15 @@ function deriveRequests(events: readonly SessionEvent[]): readonly RequestView[]
}) })
continue continue
} }
if (sourceEvent.type === 'turn/end' && sourceEvent.data.reason.kind === 'error') { if (sourceEvent.type === 'turn/end') {
const reason = sourceEvent.data.reason const lastStep = lastStepByTurn.get(sourceEvent.data.turn)
updateAssistant(ordinaryByStep.get(requestKey(sourceEvent.data.turn, sourceEvent.data.step)), { if (sourceEvent.data.reason.kind === 'error') {
status: 'error', updateAssistant(lastStep === undefined ? undefined : ordinaryByStep.get(lastStep), {
error: displayFailureMessage(reason.error), status: 'error',
}) error: displayFailureMessage(sourceEvent.data.reason.error),
})
}
lastStepByTurn.delete(sourceEvent.data.turn)
continue continue
} }
@@ -98,6 +98,8 @@ export class Session implements SessionFace {
private readonly transcript = new TranscriptAdapter() private readonly transcript = new TranscriptAdapter()
private partial: PartialAccumulator | null = null private partial: PartialAccumulator | null = null
private openCalls = new Map<string, RunningToolCall>() private openCalls = new Map<string, RunningToolCall>()
/** Last entered step per turn, folded from step/start for terminal error placement. */
private lastStepByTurn = new Map<number, number>()
/** Operational notices and interrupted-turn terminal nodes merged into the flow by seq. /** 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 * Derived from window events and rebuilt with partial/openCalls; the transcript is
* seq-monotonic, so a plain seq merge preserves event order. */ * seq-monotonic, so a plain seq merge preserves event order. */
@@ -796,6 +798,10 @@ export class Session implements SessionFace {
} }
switch (event.type) { switch (event.type) {
case 'turn/start': case 'turn/start':
this.lastStepByTurn.set(event.data.turn, 0)
return
case 'step/start':
this.lastStepByTurn.set(event.data.turn, event.data.step)
return return
case 'assistant/chunk': { case 'assistant/chunk': {
const { turn, step, chunk } = event.data const { turn, step, chunk } = event.data
@@ -826,6 +832,7 @@ export class Session implements SessionFace {
return return
} }
case 'turn/end': { case 'turn/end': {
const lastStep = this.lastStepByTurn.get(event.data.turn) ?? 0
this.turnEnds.set(event.data.turn, event.seq) this.turnEnds.set(event.data.turn, event.seq)
this.turnEndsRev++ this.turnEndsRev++
if (event.data.reason.kind === 'aborted') { if (event.data.reason.kind === 'aborted') {
@@ -841,7 +848,7 @@ export class Session implements SessionFace {
seq: event.seq, seq: event.seq,
time: event.time, time: event.time,
turn: event.data.turn, turn: event.data.turn,
step: event.data.step, step: lastStep,
message: displayFailureMessage(failure), message: displayFailureMessage(failure),
code: failure.code, code: failure.code,
}) })
@@ -882,6 +889,7 @@ export class Session implements SessionFace {
}) })
this.derivedRev++ this.derivedRev++
} }
this.lastStepByTurn.delete(event.data.turn)
return return
} }
default: default:
@@ -916,6 +924,7 @@ export class Session implements SessionFace {
private rebuildDerivedFromWindow(): void { private rebuildDerivedFromWindow(): void {
this.partial = null this.partial = null
this.openCalls.clear() this.openCalls.clear()
this.lastStepByTurn.clear()
this.callsRev++ this.callsRev++
this.derivedNodes = [] this.derivedNodes = []
this.derivedRev++ this.derivedRev++
@@ -2385,7 +2385,7 @@ export const TYPE_API: readonly TypeApiEntry[] = [
}, },
{ {
name: 'SessionEventMap', 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<string, never>;\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<string, never>;\n}',
}, },
{ {
name: 'SessionEventMetadataFilter', name: 'SessionEventMetadataFilter',
+1 -1
View File
@@ -296,7 +296,7 @@ export class ReactLoopAgent implements Agent {
} finally { } finally {
try { try {
// oxlint-disable-next-line typescript/no-non-null-assertion -- every exit assigns a turn ending // 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) { } catch (error: unknown) {
this.throwError(error) this.throwError(error)
} }
-4
View File
@@ -87,10 +87,6 @@ function validateEvent(
if (trace.openStep !== null) { if (trace.openStep !== null) {
fail(`turn/end ${event.data.turn} while step ${trace.openStep} is still open`) 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 openTurn = null
nextTurn += 1 nextTurn += 1
break break
+1 -5
View File
@@ -46,7 +46,6 @@ export const TOOL_OUTCOME_UNKNOWN = 'TOOL_OUTCOME_UNKNOWN'
export function interruptedTurnClosers(events: readonly SessionEvent[]): SessionEvent[] { export function interruptedTurnClosers(events: readonly SessionEvent[]): SessionEvent[] {
let openTurn: number | null = null let openTurn: number | null = null
let openStep: number | null = null let openStep: number | null = null
let lastStep = 0
// Reset at each turn boundary so earlier calls cannot leak into tail repair. // Reset at each turn boundary so earlier calls cannot leak into tail repair.
// Assistant blocks register calls; later tool/call events add provenance seqs. // Assistant blocks register calls; later tool/call events add provenance seqs.
const pendingCalls = new Map<CallId, { step: number; callSeq?: number }>() const pendingCalls = new Map<CallId, { step: number; callSeq?: number }>()
@@ -55,18 +54,15 @@ export function interruptedTurnClosers(events: readonly SessionEvent[]): Session
case 'turn/start': case 'turn/start':
openTurn = event.data.turn openTurn = event.data.turn
openStep = null openStep = null
lastStep = 0
pendingCalls.clear() pendingCalls.clear()
break break
case 'turn/end': case 'turn/end':
openTurn = null openTurn = null
openStep = null openStep = null
lastStep = 0
pendingCalls.clear() pendingCalls.clear()
break break
case 'step/start': case 'step/start':
openStep = event.data.step openStep = event.data.step
lastStep = event.data.step
break break
case 'step/end': case 'step/end':
pendingCalls.clear() pendingCalls.clear()
@@ -151,6 +147,6 @@ export function interruptedTurnClosers(events: readonly SessionEvent[]): Session
if (openStep !== null) { if (openStep !== null) {
closers.push({ type: 'step/end', seq: seq++, time, data: { turn: openTurn, step: openStep } }) 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 return closers
} }
+3 -3
View File
@@ -196,14 +196,14 @@ export interface SessionEventMap {
*/ */
'turn/start': { turn: number } 'turn/start': { turn: number }
/** /**
* Closes turn `turn` after `step`, the last entered step (`0` when none), * Closes turn `turn` with the {@link TurnEndReason} that ended it. A turn
* with the {@link TurnEndReason} that ended it. The loop does not await a * 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 * flush at turn boundaries: `dsh-session-checkpoint-policy` owns the
* per-request durability checkpoint, and consumers that read storage after * per-request durability checkpoint, and consumers that read storage after
* `whenIdle()` flush themselves. Success commits the turn; rejection is * `whenIdle()` flush themselves. Success commits the turn; rejection is
* reported live and does not prevent later work. * 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. */ /** Opens step `step` of turn `turn` — one model call plus the tool executions it requested. */
'step/start': { turn: number; step: number } 'step/start': { turn: number; step: number }
/** Closes step `step` of turn `turn`. */ /** Closes step `step` of turn `turn`. */
+6 -5
View File
@@ -320,7 +320,9 @@ export function apply(ctx: Context): void {
return return
} }
if (event.data.reason.kind !== 'aborted') 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) else disarm(state)
return return
default: default:
@@ -372,10 +374,9 @@ export function apply(ctx: Context): void {
decision = await next() decision = await next()
} catch (error: unknown) { } catch (error: unknown) {
if (signal.aborted) throw error if (signal.aborted) throw error
// A throwing downstream hook drops the whole step proposal: the loop // A throwing downstream hook drops the whole step proposal. Clear the
// returns to idle without a turn, so a still-queued reservation would // reservation before the balanced no-step turn returns to idle so the
// starve every later drive pass. Clear it and let the driver // next drive pass can reschedule the round.
// reschedule the round.
state.attempt = undefined state.attempt = undefined
requestDrive(state) requestDrive(state)
throw error throw error
@@ -71,7 +71,7 @@ export interface PersistenceBackend<TornMarker = unknown> {
* region strictly below `fromSeq` is limited to seq contiguity — the * region strictly below `fromSeq` is limited to seq contiguity — the
* service contract scopes this read to the suffix — unless that suffix * service contract scopes this read to the suffix — unless that suffix
* contains a supported legacy shape whose normalization needs earlier * 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. * to the complete stored prefix.
* @param id - persisted session id to resolve. * @param id - persisted session id to resolve.
* @param fromSeq - first event seq to include (non-negative safe integer, * @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 data = asRecord(event.data)
const legacySteeringType: string = 'steering/message' const legacySteeringType: string = 'steering/message'
if (event.type === legacySteeringType) return true if (event.type === legacySteeringType) return true
if (event.type === 'turn/end' && data !== undefined && !Object.hasOwn(data, 'step')) return true
if (data === undefined) return false if (data === undefined) return false
switch (event.type) { switch (event.type) {
case 'user/message': case 'user/message':
@@ -270,11 +269,11 @@ function migrateLegacyTurnStartEvent(event: SessionEvent, id: SessionId): Sessio
return { ...event, data: { turn: data['turn'] } } as SessionEvent return { ...event, data: { turn: data['turn'] } } as SessionEvent
} }
/** Upgrade the turn boundary emitted immediately before the loop refactor. */ /** Upgrade an obsolete turn ending while preserving the latest-master envelope. */
function migrateLegacyTurnEndEvent(event: SessionEvent, id: SessionId, lastStep: number): SessionEvent { function migrateLegacyTurnEndEvent(event: SessionEvent, id: SessionId): SessionEvent {
if (event.type !== 'turn/end') return event if (event.type !== 'turn/end') return event
const data = asRecord(event.data) const data = asRecord(event.data)
if (data === undefined || Object.hasOwn(data, 'step')) return event if (data === undefined) return event
const malformed = (): never => { const malformed = (): never => {
throw new Error(`session "${id}" contains malformed pre-react-loop turn/end at seq ${event.seq}`) 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']) || !hasOnlyKeys(data, ['turn', 'reason'])
|| reason === undefined || typeof reason['kind'] !== 'string') return malformed() || reason === undefined || typeof reason['kind'] !== 'string') return malformed()
let currentReason: Record<string, unknown> let currentReason: Record<string, unknown> | undefined
switch (reason['kind']) { switch (reason['kind']) {
case 'completed': case 'completed':
case 'blocked':
case 'max-tokens': case 'max-tokens':
case 'interrupted': case 'interrupted':
if (!hasOnlyKeys(reason, ['kind'])) return malformed() if (!hasOnlyKeys(reason, ['kind'])) return malformed()
currentReason = { kind: reason['kind'] } return event
break
case 'aborted': case 'aborted':
if (Object.hasOwn(reason, 'reason')) return event
if (!hasOnlyKeys(reason, ['kind'])) return malformed() if (!hasOnlyKeys(reason, ['kind'])) return malformed()
currentReason = { kind: 'aborted', reason: { kind: 'legacy' } } currentReason = { kind: 'aborted', reason: { kind: 'legacy' } }
break break
@@ -300,6 +300,7 @@ function migrateLegacyTurnEndEvent(event: SessionEvent, id: SessionId, lastStep:
currentReason = { kind: 'aborted', reason: { kind: 'disposed' } } currentReason = { kind: 'aborted', reason: { kind: 'disposed' } }
break break
case 'error': { case 'error': {
if (Object.hasOwn(reason, 'error')) return event
if (!Number.isSafeInteger(reason['step']) || (reason['step'] as number) < 0) return malformed() if (!Number.isSafeInteger(reason['step']) || (reason['step'] as number) < 0) return malformed()
const failure = asRecord(reason['failure']) const failure = asRecord(reason['failure'])
if (failure !== undefined && hasOnlyKeys(reason, ['kind', 'step', 'failure']) if (failure !== undefined && hasOnlyKeys(reason, ['kind', 'step', 'failure'])
@@ -327,14 +328,13 @@ function migrateLegacyTurnEndEvent(event: SessionEvent, id: SessionId, lastStep:
break break
} }
default: default:
return malformed() return event
} }
return { return {
...event, ...event,
data: { data: {
...data, ...data,
step: lastStep,
reason: currentReason, reason: currentReason,
}, },
} as SessionEvent } as SessionEvent
@@ -431,16 +431,9 @@ function eventMessageId(event: SessionEvent): PersistedMessageId | undefined {
function snapshotStoredEvents(events: readonly SessionEvent[], id: SessionId): SessionEvent[] { function snapshotStoredEvents(events: readonly SessionEvent[], id: SessionId): SessionEvent[] {
assertSupportedEvents(events, id) assertSupportedEvents(events, id)
const messageIds = new Map<number, PersistedMessageId>() const messageIds = new Map<number, PersistedMessageId>()
const lastSteps = new Map<number, number>()
return events.map((event) => { 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 migratedStart = migrateLegacyTurnStartEvent(event, id)
const migratedTurn = migrateLegacyTurnEndEvent(migratedStart, id, lastStep) const migratedTurn = migrateLegacyTurnEndEvent(migratedStart, id)
const migratedSteering = migrateLegacySteeringEvent(migratedTurn, id) const migratedSteering = migrateLegacySteeringEvent(migratedTurn, id)
const snapshot = snapshotSessionEvent(migrateLegacyMessageEvent(migratedSteering, id, messageIds)) const snapshot = snapshotSessionEvent(migrateLegacyMessageEvent(migratedSteering, id, messageIds))
const messageId = eventMessageId(snapshot) const messageId = eventMessageId(snapshot)