From 218bb7f6452793de50c53e39cdbb5e8db463a4b1 Mon Sep 17 00:00:00 2001 From: Yichen Jiang Date: Sun, 30 Aug 2026 12:54:52 +0800 Subject: [PATCH] feat(session-controller): loadThrough deep history paging Session.loadThrough(seq) loops the existing prepend pager (200-message pages) until the window covers the target seq, with a shared low-water retarget for repeated calls, a no-progress guard against empty pages still claiming history, and loadOlder's fail-soft error posture. Busy state rides the existing loadingOlder snapshot bit. --- .../src/client/contract/session.ts | 9 ++ .../src/client/sessions/session.ts | 42 +++++++++ .../tests/session.client.spec.ts | 90 ++++++++++++++++++- .../conversation-registry.client.spec.ts | 1 + .../client-runtime/src/sessions.ts | 8 ++ 5 files changed, 149 insertions(+), 1 deletion(-) diff --git a/packages/api/session-controller/src/client/contract/session.ts b/packages/api/session-controller/src/client/contract/session.ts index 02214bb872..71cc72e0ab 100644 --- a/packages/api/session-controller/src/client/contract/session.ts +++ b/packages/api/session-controller/src/client/contract/session.ts @@ -119,6 +119,15 @@ export interface ISession { * @returns completion; failures land in snapshot.openState/loadingOlder. */ loadOlder(): Promise + /** + * Page history backwards until the window covers `seq` (inclusive) — the + * turn-jump loader. Repeated calls while a jump is paging lower its shared + * target and return the in-flight completion; `snapshot.loadingOlder` is + * the busy signal for the whole jump. + * @param seq - durable event seq the window must reach (a turn's `turn/start` seq). + * @returns completion once covered, exhausted, superseded, or failed soft. + */ + loadThrough(seq: number): Promise /** * Execute one slash-command line against this session's agent — pure * admission semantics (the host executor durably logs the lifecycle). diff --git a/packages/api/session-controller/src/client/sessions/session.ts b/packages/api/session-controller/src/client/sessions/session.ts index 25d5e61082..294f33f3aa 100644 --- a/packages/api/session-controller/src/client/sessions/session.ts +++ b/packages/api/session-controller/src/client/sessions/session.ts @@ -38,6 +38,9 @@ import { SessionQueueMirror } from './queue-mirror.ts' /** Messages requested per history page. */ export const PAGE_MESSAGES = 50 +/** Messages requested per page while a turn jump loops backwards (fewer, larger round trips). */ +export const JUMP_PAGE_MESSAGES = 200 + /** Manager-owned observers of a Session object's local state edges. */ export interface SessionOptions { /** Catalog-discovered address selecting non-activating subagent transport. */ @@ -78,6 +81,10 @@ export class Session implements SessionFace { * passes drop all writes once the generation moves on. */ private openGeneration = 0 private loadingOlder = false + /** Shared low-water target of the running jump loop; null when no jump is paging. */ + private jumpTargetSeq: number | null = null + /** The running jump loop's completion, shared by retargeting callers. */ + private jumpPromise: Promise | null = null /** Authoritative stream-only inbox snapshot; pending work never hits history. */ private readonly queueMirror = new SessionQueueMirror() private running = false @@ -369,6 +376,41 @@ export class Session implements SessionFace { } } + /** Jump loader: page backwards until the window covers seq (see ISession.loadThrough). */ + loadThrough(seq: number): Promise { + if (this.openState !== 'open' || !this.hasMore || this.baseSeq <= seq) return Promise.resolve() + this.jumpTargetSeq = Math.min(this.jumpTargetSeq ?? seq, seq) + if (this.jumpPromise !== null) return this.jumpPromise + // A plain single-page pull owns the busy flag; the jump does not queue + // behind it (the caller may retry once it settles). + if (this.loadingOlder) return Promise.resolve() + this.loadingOlder = true + this.notifier.markDirty() + this.jumpPromise = (async () => { + try { + while (this.hasMore && this.jumpTargetSeq !== null && this.baseSeq > this.jumpTargetSeq) { + const events = this.events + if (events === undefined) return + const before = this.baseSeq + await events.prepend({ beforeSeq: this.baseSeq, maxMessages: JUMP_PAGE_MESSAGES }) + // No-progress guard: an empty or dropped page that still claims more + // history must end the loop, not spin it. + if (this.baseSeq >= before) return + } + } catch (error) { + if (!isRemoteFailure(error)) { + console.error('[session-controller] loadThrough failed:', error) + } + } finally { + this.jumpTargetSeq = null + this.jumpPromise = null + this.loadingOlder = false + this.notifier.markDirty() + } + })() + return this.jumpPromise + } + /** Rebuild an opened history source after address replacement. * Invalidates any in-flight open first; queue state belongs to the independently * reconnecting control stream and remains untouched. */ diff --git a/packages/api/session-controller/tests/session.client.spec.ts b/packages/api/session-controller/tests/session.client.spec.ts index 319c315873..2bee7d33b4 100644 --- a/packages/api/session-controller/tests/session.client.spec.ts +++ b/packages/api/session-controller/tests/session.client.spec.ts @@ -5,7 +5,7 @@ import type { SessionEvent } from '@deepseek-ai/dsh-session/types' import type { SessionId } from '@deepseek-ai/dsh-api-remotes/client' import { RemoteStreamCarrierError } from '@deepseek-ai/dsh-api-gateway/client' import { RemoteError } from '@deepseek-ai/dsh-typert-protocol' -import { Session, type SessionOptions } from '../src/client/sessions/session.ts' +import { JUMP_PAGE_MESSAGES, Session, type SessionOptions } from '../src/client/sessions/session.ts' import { FakeApiClient, deferred, err, fakeRemote, ok } from './fake-api.client.ts' import { entries, ev, historyValue, plainTurn } from './event-script.client.ts' @@ -254,6 +254,94 @@ describe('paging', () => { } }) + it('loadThrough pages repeatedly until the window covers the target seq', async () => { + const oldest = plainTurn(0, 0, '最旧问', '最旧答') + const middle = plainTurn(6, 1, '中问', '中答') + const newest = plainTurn(12, 2, '新问', '新答') + const { api, session } = makeSession() + api.onHistory = (payload) => { + if (payload.beforeSeq === undefined) return histResponse(newest, true) + return payload.beforeSeq === 12 ? histResponse(middle, true) : histResponse(oldest, false) + } + await session.open() + + const gate = deferred>>() + api.onHistory = (payload) => { + api.onHistory = payload2 => payload2.beforeSeq === 12 ? histResponse(middle, true) : histResponse(oldest, false) + void payload + return gate.promise + } + const jump = session.loadThrough(0) + expect(session.getSnapshot().loadingOlder).toBe(true) + gate.resolve(ok(historyValue(middle, true))) + await jump + const snapshot = session.getSnapshot() + expect(snapshot.loadingOlder).toBe(false) + expect(eventSeqs(session)).toEqual([...oldest, ...middle, ...newest].map(event => event.seq)) + expect(api.callsOf('session.history')).toMatchObject([ + { beforeSeq: 12, maxMessages: JUMP_PAGE_MESSAGES }, + { beforeSeq: 6, maxMessages: JUMP_PAGE_MESSAGES }, + ]) + }) + + it('loadThrough is a no-op when the window already covers the target or the session is not open', async () => { + const { api, session } = makeSession() + await session.loadThrough(0) // cold: no-op + expect(api.calls).toEqual([]) + api.onHistory = () => histResponse(plainTurn(6, 1, 'x', 'y'), true) + await session.open() + const calls = api.calls.length + await session.loadThrough(6) // baseSeq is already 6 + await session.loadThrough(9) // inside the window + expect(api.calls.length).toBe(calls) + }) + + it('loadThrough retargets a running jump to the lowest requested seq and shares its completion', async () => { + const oldest = plainTurn(0, 0, 'a', 'b') + const middle = plainTurn(6, 1, 'c', 'd') + const { api, session } = makeSession() + api.onHistory = () => histResponse(plainTurn(12, 2, 'e', 'f'), true) + await session.open() + + const gate = deferred>>() + api.onHistory = () => { + api.onHistory = () => histResponse(oldest, false) + return gate.promise + } + const first = session.loadThrough(6) + const second = session.loadThrough(0) + gate.resolve(ok(historyValue(middle, true))) + await Promise.all([first, second]) + expect(eventSeqs(session)).toEqual([...oldest, ...middle].map(event => event.seq).concat([12, 13, 14, 15, 16, 17])) + expect(api.callsOf('session.history')).toHaveLength(2) + }) + + it('loadThrough stops on a page that makes no progress instead of looping', async () => { + const { api, session } = makeSession() + api.onHistory = payload => payload.beforeSeq === undefined + ? histResponse(plainTurn(12, 2, 'x', 'y'), true) + : histResponse([], true) // empty page still claiming more history + await session.open() + await session.loadThrough(0) + expect(session.getSnapshot().loadingOlder).toBe(false) + expect(api.callsOf('session.history')).toHaveLength(1) + }) + + it('loadThrough fails soft on a thrown page and clears its busy state', async () => { + const { api, session } = makeSession() + api.onHistory = () => histResponse(plainTurn(12, 2, 'x', 'y'), true) + await session.open() + api.onHistory = () => Promise.reject(new Error('page wire down')) + const errorSpy = vi.spyOn(console, 'error').mockImplementation(() => undefined) + try { + await session.loadThrough(0) + expect(errorSpy).toHaveBeenCalled() + expect(session.getSnapshot().loadingOlder).toBe(false) + } finally { + errorSpy.mockRestore() + } + }) + it('ignores loadOlder while one is in flight (single request)', async () => { const { api, session } = makeSession() api.onHistory = () => histResponse(plainTurn(6, 1, 'x', 'y'), true) diff --git a/packages/client/ui-conversation/tests/conversation-registry.client.spec.ts b/packages/client/ui-conversation/tests/conversation-registry.client.spec.ts index 63b2c22146..ba52e36ba6 100644 --- a/packages/client/ui-conversation/tests/conversation-registry.client.spec.ts +++ b/packages/client/ui-conversation/tests/conversation-registry.client.spec.ts @@ -51,6 +51,7 @@ function fakeSession(): SessionFace { cancel: () => Promise.reject(new Error('unused fake Session operation')), rename: () => Promise.reject(new Error('unused fake Session operation')), loadOlder: () => Promise.reject(new Error('unused fake Session operation')), + loadThrough: () => Promise.reject(new Error('unused fake Session operation')), command: () => Promise.reject(new Error('unused fake Session operation')), } } diff --git a/packages/test-support/client-runtime/src/sessions.ts b/packages/test-support/client-runtime/src/sessions.ts index da67718820..419e155af7 100644 --- a/packages/test-support/client-runtime/src/sessions.ts +++ b/packages/test-support/client-runtime/src/sessions.ts @@ -152,6 +152,14 @@ export class FixtureSession implements SessionFace { throw new Error(`test session "${this.sessionId}": loadOlder is not stubbed — supply it on the fixture's session face`) } + /** + * Fail-loud stub; supply `loadThrough` on the fixture's session face to exercise it. + * @returns never — always throws. + */ + loadThrough(): never { + throw new Error(`test session "${this.sessionId}": loadThrough is not stubbed — supply it on the fixture's session face`) + } + /** * Fail-loud stub; supply `rename` on the fixture's session face to exercise it. * @returns never — always throws.