diff --git a/packages/client/runtime/src/client/sessions/session.ts b/packages/client/runtime/src/client/sessions/session.ts index 49035736dc..6708202276 100644 --- a/packages/client/runtime/src/client/sessions/session.ts +++ b/packages/client/runtime/src/client/sessions/session.ts @@ -82,9 +82,8 @@ export class Session implements ObservableSnapshot { private openState: OpenState = 'cold' private openError: RpcError | null = null private openPromise: Promise | null = null - /** Bumped by resync to invalidate an in-flight doOpen: a reconnect must rebuild, never adopt - * a pre-disconnect open whose history request is already doomed (audit S4). Stale doOpen - * passes drop all writes once the generation moves on. */ + /** Bumped at disconnect and resync to invalidate in-flight history work: a reconnect must + * rebuild, never adopt a pre-disconnect response (audit S4). */ private openGeneration = 0 private loadingOlder = false private readonly foldAdapter = new FoldAdapter() @@ -270,12 +269,14 @@ export class Session implements ObservableSnapshot { /** Page up: pull one earlier page with the window's first seq as beforeSeq and prepend (§D.2). */ async loadOlder(): Promise { if (this.openState !== 'open' || !this.hasMore || this.loadingOlder) return + const generation = this.openGeneration this.loadingOlder = true this.notifier.markDirty() try { const { result } = await this.api.sessions.history({ sessionId: this.sessionId, beforeSeq: this.baseSeq, maxMessages: PAGE_MESSAGES, }) + if (generation !== this.openGeneration) return if (!result.ok) return // keep the window as-is; do not overwrite openError (open already succeeded) const older = result.value.events if (older.length === 0) { @@ -299,8 +300,10 @@ export class Session implements ObservableSnapshot { } catch (error) { console.error('[web-runtime] loadOlder failed:', error) } finally { - this.loadingOlder = false - this.notifier.markDirty() + if (generation === this.openGeneration) { + this.loadingOlder = false + this.notifier.markDirty() + } } } @@ -321,6 +324,8 @@ export class Session implements ObservableSnapshot { this.events = [] this.views = [] this.baseSeq = 0 + this.loadingOlder = false + this.stitching = false // Superseded, not settled: the baseline replay re-sends still-pending requested frames verbatim // (same rpcId), re-minting fresh waits; a stale reference's respond() still reaches the host. this.pending.clear() @@ -483,6 +488,7 @@ export class Session implements ObservableSnapshot { /** Connection-loss boundary: clear values that are not replayed before the next stream starts. */ handleReconnecting(): void { + this.openGeneration++ if (this.metrics === null && this.contextWindow === undefined) return this.metrics = null this.contextWindow = undefined @@ -649,7 +655,7 @@ export class Session implements ObservableSnapshot { } catch (error) { console.error('[web-runtime] gap repair failed:', error) } finally { - this.stitching = false + if (generation === this.openGeneration) this.stitching = false } } diff --git a/packages/client/runtime/tests/session.spec.ts b/packages/client/runtime/tests/session.spec.ts index cbce270597..d879362990 100644 --- a/packages/client/runtime/tests/session.spec.ts +++ b/packages/client/runtime/tests/session.spec.ts @@ -392,6 +392,27 @@ describe('paging', () => { await Promise.all([first, second]) expect(api.callsOf('session.history')).toHaveLength(2) // open + one page, not two }) + + it('drops an older page from the disconnected generation', async () => { + const { api, session } = makeSession() + api.onHistory = () => histResponse(plainTurn(6, 1, '新问', '新答'), true) + await session.open() + const stale = deferred>>() + api.onHistory = () => stale.promise + const loading = session.loadOlder() + + session.handleReconnecting() + stale.resolve(ok({ + events: entries(plainTurn(0, 0, '旧问', '旧答')) as never[], + hasMore: false, + })) + await loading + expect(session.getSnapshot().nodes.map(node => node.seq)).toEqual([7, 9]) + + api.onHistory = () => histResponse(plainTurn(12, 2, '重连问', '重连答')) + await session.resync() + expect(session.getSnapshot().nodes.map(node => node.seq)).toEqual([13, 15]) + }) }) describe('prompt and cancel errors', () => { @@ -733,22 +754,49 @@ describe('remaining branches', () => { expect(session.getSnapshot().openState).toBe('open') }) - it('drops a gap repair superseded by a full resync while its pull was in flight', async () => { + it('drops a stale gap repair without clearing a newer generation repair', async () => { const { api, session } = makeSession() api.onHistory = () => histResponse(plainTurn(0, 0, 'a', 'b')) await session.open() - const repairPull = deferred>>() - api.onHistory = () => repairPull.promise + const staleRepair = deferred>>() + api.onHistory = () => staleRepair.promise session.handleMuxEnvelope('r' as never, { type: 'session/event', sessionId: SID, event: ev.user(9, '洞') }) // starts repairGap + session.handleReconnecting() api.onHistory = () => histResponse(plainTurn(6, 1, 'c', 'd')) - const resynced = session.resync() // bumps the generation - repairPull.resolve(ok({ + await session.resync() + + const freshRepair = deferred>>() + let freshRepairCalls = 0 + api.onHistory = () => { + freshRepairCalls++ + return freshRepair.promise + } + session.handleMuxEnvelope('fresh-gap' as never, { + type: 'session/event', + sessionId: SID, + event: ev.user(15, '新洞'), + }) + expect(freshRepairCalls).toBe(1) + + staleRepair.resolve(ok({ events: entries(plainTurn(0, 0, '旧', '页')) as never[], hasMore: false, - modelTarget: { provider: 'deepseek', model: 'stale' }, })) // repair result: stale, dropped - await resynced - expect(session.getSnapshot().nodes.map(n => n.seq)).toEqual([7, 9]) + await Promise.resolve() + session.handleMuxEnvelope('fresh-buffer' as never, { + type: 'session/event', + sessionId: SID, + event: ev.user(16, '继续缓存'), + }) + expect(freshRepairCalls).toBe(1) // stale finally did not clear the newer stitching owner + + freshRepair.resolve(ok({ + events: entries([...plainTurn(6, 1, 'c', 'd'), ...plainTurn(12, 2, 'e', 'f')]) as never[], + hasMore: false, + })) + await vi.waitFor(() => { + expect(session.getSnapshot().nodes.map(n => n.seq)).toEqual([7, 9, 13, 15]) + }) }) it('successful cancel leaves no promptError; tool/result for an unknown callId is a no-op', async () => { @@ -813,6 +861,53 @@ describe('remaining branches', () => { }) describe('resync', () => { + it('fences pre-disconnect history behind a fresh mux metrics baseline', async () => { + const { api, session } = makeSession() + const stale = deferred>>() + api.onHistory = () => stale.promise + const opening = session.open() + const oldLiveMetrics = metrics(8, 10) + session.handleMuxEnvelope('old-metrics' as never, { + type: 'session/metrics', + sessionId: SID, + metrics: oldLiveMetrics, + }) + session.handleMuxEnvelope('old-capacity' as never, { + type: 'session/model-request', + sessionId: SID, + turn: 1, + step: 1, + provider: 'test', + model: 'old', + contextWindow: 128_000, + }) + + session.handleReconnecting() + expect(session.getSnapshot().metrics).toBeNull() + expect(session.getSnapshot().modelRequestContextWindow).toBeUndefined() + const freshMetrics = metrics(0, 1, { contextTokens: 20 }) + session.handleMuxEnvelope('fresh-metrics' as never, { + type: 'session/metrics', + sessionId: SID, + metrics: freshMetrics, + }) + + stale.resolve(ok({ + events: entries(plainTurn(0, 0, '旧问', '旧答')) as never[], + hasMore: false, + metrics: metrics(99, 99, { contextTokens: 999 }), + })) + await opening + expect(session.getSnapshot().nodes).toEqual([]) + expect(session.getSnapshot().metrics).toBe(freshMetrics) + + api.onHistory = () => histResponse(plainTurn(6, 1, '新问', '新答')) + await session.resync() + expect(session.getSnapshot().openState).toBe('open') + expect(session.getSnapshot().nodes.map(node => node.seq)).toEqual([7, 9]) + expect(session.getSnapshot().metrics).toBe(freshMetrics) + }) + it('preserves fresh-generation metrics that arrive before a failing history refresh', async () => { const { api, session } = makeSession() const oldMetrics = metrics(8, 10)