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.
This commit is contained in:
parent
7e2eacb1fe
commit
218bb7f645
5 changed files with 149 additions and 1 deletions
|
|
@ -119,6 +119,15 @@ export interface ISession {
|
|||
* @returns completion; failures land in snapshot.openState/loadingOlder.
|
||||
*/
|
||||
loadOlder(): Promise<void>
|
||||
/**
|
||||
* 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<void>
|
||||
/**
|
||||
* Execute one slash-command line against this session's agent — pure
|
||||
* admission semantics (the host executor durably logs the lifecycle).
|
||||
|
|
|
|||
|
|
@ -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<void> | 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<void> {
|
||||
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. */
|
||||
|
|
|
|||
|
|
@ -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<Awaited<ReturnType<FakeApiClient['onHistory']>>>()
|
||||
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<Awaited<ReturnType<FakeApiClient['onHistory']>>>()
|
||||
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)
|
||||
|
|
|
|||
|
|
@ -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')),
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue