1
/** Cold Session history pagination and live-event source. */3
import type { Context } from '@deepseek-ai/cordis'4
import { Deque } from '@deepseek-ai/dsh-deque'5
import type { AssistantStreamFrame } from '@deepseek-ai/dsh-agent'6
import {7
isAppendSurfaceEvent,8
SessionLogOffset,9
SessionSeq,10
} from '@deepseek-ai/dsh-session'11
import type {12
SessionEvent,13
SessionHeader,14
SessionId,15
SessionLogOffset as SessionLogOffsetType,16
SessionSeqCursor,17
} from '@deepseek-ai/dsh-session'18
import { SessionQueryError, type SessionObservation } from '@deepseek-ai/dsh-session-query'19
import type {} from '@deepseek-ai/dsh-subagent'20
import { RemoteError } from '@deepseek-ai/dsh-typert-protocol'21
import type { JsonValue } from '@deepseek-ai/dsh-util-values'22
import type {23
SessionAddress,24
SessionAssistantStreamFrame,25
SessionEventEntry,26
SessionFollowRequest,27
SessionFollowFrame,28
SessionHistoryRecord,29
SessionPage,30
SessionPageRequest,31
SessionProjectionBaseline,32
SessionProjectionValues,33
SessionWireHeader,34
SessionWireEvent,35
} from './types.ts'36
import { SessionAssistantStreamAccumulator } from './assistant-stream.ts'38
const DEFAULT_MAX_MESSAGES = 5039
const MESSAGE_TYPES = new Set(['user/message', 'assistant/message'])41
/** Implements cold-safe history operations delegated by the Session Controller. */42
export class SessionHistoryController {43
private readonly closeFollowers = new Set<() => void>()44
private readonly assistantStreams = new Map<SessionId, SessionAssistantStreamAccumulator>()46
/**47
* @param ctx - Host context carrying Session query and projection services.48
* @param promote - starts ordinary Session activation after snapshot delivery.49
*/50
constructor(51
private readonly ctx: Context,52
private readonly promote: (observation: SessionObservation) => void,53
) {54
ctx.on('agent/assistant-stream', ({ agent, frame }) => {55
let stream = this.assistantStreams.get(agent.session.id)56
if (stream === undefined) {57
stream = new SessionAssistantStreamAccumulator()58
this.assistantStreams.set(agent.session.id, stream)59
}60
stream.accept(frame, cursorBeforeNext(agent.session.seq))61
}, { global: true })62
ctx.on('agent/disposed', ({ agent }) => {63
this.assistantStreams.delete(agent.session.id)64
}, { global: true })65
ctx.effect(() => () => {66
for (const close of this.closeFollowers) close()67
this.closeFollowers.clear()68
}, 'session-controller.history')69
}71
/**72
* Read one message-aligned history page without activating an Agent.73
* @param request - durable address and backwards-page cursor.74
* @param signal - caller cancellation for persistence reads.75
* @returns a contiguous event page.76
*/77
async page(request: SessionPageRequest, signal: AbortSignal): Promise<SessionPage> {78
validatePageRequest(request)79
const throughSeq: SessionSeqCursor = request.throughSeq === -180
? -181
: SessionSeq(request.throughSeq)82
const beforeSeq = request.beforeSeq === undefined83
? undefined84
: SessionLogOffset(request.beforeSeq)85
using source = await this.sourceFor(request.address, signal, false)86
signal.throwIfAborted()87
const sourceLog = source.events88
const sourceCursor: SessionSeqCursor = sourceLog.at(-1)?.seq ?? -189
if (throughSeq > sourceCursor) {90
throw new RemoteError(91
'gateway/bad-request',92
`session page through seq ${String(throughSeq)} is past cursor ${String(sourceCursor)}`,93
{},94
)95
}96
/* v8 ignore next -- Session and persistence validation guarantee a dense zero-based event prefix. */97
if (throughSeq >= 0 && sourceLog[throughSeq]?.seq !== throughSeq) {98
throw new RemoteError('gateway/internal', `session log does not contain through seq ${String(throughSeq)}`, {})99
}100
const page = paginate(101
sourceLog,102
beforeSeq,103
request.maxMessages ?? DEFAULT_MAX_MESSAGES,104
throughSeq,105
request.turnWindow,106
)107
const records = pageRecords(page.events)108
return {109
records,110
hasMore: page.hasMore,111
}112
}114
/**115
* Follow events appended after an initial cursor on one durable address.116
* @param request - durable address and last committed sequence already held by the caller.117
* @param signal - stream cancellation owned by the Remote carrier.118
* @returns a complete opening snapshot followed by gap-free durable events and opted-in assistant frames.119
*/120
async *follow(request: SessionFollowRequest, signal: AbortSignal): AsyncIterable<SessionFollowFrame> {121
validateHistoryWindow(request)122
const { address } = request123
const target = addressId(address)124
const buffered = new Deque<125
| { readonly type: 'event'; readonly event: SessionEvent }126
| {127
readonly type: 'assistant-stream'128
readonly frame: SessionAssistantStreamFrame129
readonly ordinal: number130
}131
>()132
let snapshotCursor: SessionSeqCursor | undefined133
let assistantStreamOrdinal = 0134
let wake: (() => void) | undefined135
const notify = (): void => {136
const resume = wake137
wake = undefined138
resume?.()139
}140
const follower = { closed: false }141
const close = (): void => {142
follower.closed = true143
notify()144
}145
this.closeFollowers.add(close)146
const disposeEvent = this.ctx.on('session/event', (session, event) => {147
if (session.id !== target) return148
buffered.pushBack({ type: 'event', event })149
notify()150
}, { global: true })151
const disposeCreated = this.ctx.on('session/created', (session) => {152
if (session.id !== target) return153
// Constructor seed events have no session/event notification. Normally154
// only the end-seed suffix is new; if persistence advanced after the155
// opening observation, replay everything beyond that snapshot cursor.156
// oxlint-disable-next-line typescript/no-deprecated -- Existing Session history read; migration deferred.157
const suffix = session.snapshotEvents(snapshotCursor === undefined158
? session.firstLiveSeq159
: SessionLogOffset(snapshotCursor + 1))160
for (let index = suffix.length - 1; index >= 0; index -= 1) {161
buffered.pushFront({ type: 'event', event: suffix[index] as SessionEvent })162
}163
notify()164
}, { global: true })165
const disposeAssistantStream = request.assistantStream !== true166
? undefined167
: this.ctx.on('agent/assistant-stream', ({ agent, frame }) => {168
if (agent.session.id !== target) return169
buffered.pushBack({170
type: 'assistant-stream',171
frame: wireAssistantStreamFrame(frame, cursorBeforeNext(agent.session.seq)),172
ordinal: ++assistantStreamOrdinal,173
})174
notify()175
}, { global: true })176
const onAbort = (): void => { notify() }177
signal.addEventListener('abort', onAbort, { once: true })178
try {179
using source = await this.sourceFor(address, signal, true)180
const events = source.events181
signal.throwIfAborted()182
const cursor = source.cursor183
snapshotCursor = cursor184
const page = paginate(events, undefined, request.maxMessages ?? DEFAULT_MAX_MESSAGES, cursor, request.turnWindow)185
const assistantStream = request.assistantStream === true186
? this.assistantStreams.get(target)?.snapshot() ?? { revision: 0 }187
: undefined188
// The accumulator snapshot and this watermark are synchronous. Frames189
// through the cut are represented or superseded by that baseline,190
// including larger revisions from a retired Agent; later revision191
// resets reach Client continuity validation.192
const assistantStreamOrdinalCut = assistantStreamOrdinal193
yield {194
type: 'snapshot',195
header: wireHeader(source.header),196
cursor,197
records: pageRecords(page.events),198
hasMore: page.hasMore,199
projections: source.projections === undefined200
? { asOfSeq: cursor, values: {} }201
: projectionBlock(source.projections),202
...assistantStream === undefined ? {} : { assistantStream },203
}204
if (address.kind === 'session' && source.source === 'prepared') {205
const promotion = source.retain()206
try {207
this.promote(promotion)208
} catch (error: unknown) {209
promotion[Symbol.dispose]()210
throw error211
}212
}213
let nextOffset = SessionLogOffset(cursor + 1)214
while (!follower.closed && !signal.aborted) {215
const item = buffered.popFront()216
if (item === undefined) {217
await new Promise<void>((resolve) => { wake = resolve })218
continue219
}220
if (item.type === 'assistant-stream') {221
if (item.ordinal > assistantStreamOrdinalCut) {222
yield { type: 'assistant-stream', frame: item.frame }223
}224
continue225
}226
const expectedSeq = SessionSeq(nextOffset)227
if (item.event.seq < expectedSeq) continue228
if (item.event.seq !== expectedSeq) {229
throw new RemoteError('gateway/internal', `session event stream skipped seq ${String(expectedSeq)}`, {})230
}231
nextOffset = SessionLogOffset(nextOffset + 1)232
yield entryFor(item.event)233
}234
} finally {235
this.closeFollowers.delete(close)236
signal.removeEventListener('abort', onAbort)237
disposeCreated()238
disposeEvent()239
disposeAssistantStream?.()240
}241
}243
private async sourceFor(244
address: SessionAddress,245
signal: AbortSignal,246
withProjections: boolean,247
): Promise<SessionObservation> {248
const sessionId = addressId(address)249
try {250
const observation = await this.ctx.sessionQuery.observeSession(sessionId, {251
signal,252
projectionMode: withProjections || address.kind === 'subagent' ? 'all' : 'none',253
})254
if (observation.header.cwd === undefined) {255
observation[Symbol.dispose]()256
rejectNotFound(address)257
}258
try {259
validateAddress(260
address,261
observation.header,262
observation.inheritedEventCount,263
observation.projections,264
)265
} catch (error: unknown) {266
observation[Symbol.dispose]()267
throw error268
}269
return observation270
} catch (error: unknown) {271
if (error instanceof SessionQueryError272
&& error.code === 'SESSION_QUERY_SESSION_NOT_FOUND') rejectNotFound(address)273
throw error274
}275
}277
}279
function cursorBeforeNext(nextSeq: SessionLogOffsetType): SessionSeqCursor {280
return nextSeq === 0 ? -1 : SessionSeq(nextSeq - 1)281
}283
function wireAssistantStreamFrame(284
frame: AssistantStreamFrame,285
durableCursor: SessionSeqCursor,286
): SessionAssistantStreamFrame {287
if (frame.type === 'start') return { ...frame, startedAfterSeq: durableCursor }288
if (frame.type === 'end') return frame289
return {290
...frame,291
chunk: frame.chunk as JsonValue,292
}293
}295
function projectionBlock(296
snapshot: NonNullable<SessionObservation['projections']>,297
): SessionProjectionBaseline {298
return {299
asOfSeq: snapshot.asOfSeq,300
// Projection definitions validate whole JSON values before snapshot publication.301
values: snapshot.values as SessionProjectionValues,302
}303
}305
function validatePageRequest(request: SessionPageRequest): void {306
if (!Number.isSafeInteger(request.throughSeq)307
|| request.throughSeq < -1308
|| Object.is(request.throughSeq, -0)) {309
throw new RemoteError('gateway/bad-request', 'throughSeq must be an integer greater than or equal to -1', {})310
}311
if (request.beforeSeq !== undefined312
&& (!Number.isSafeInteger(request.beforeSeq)313
|| request.beforeSeq < 0314
|| Object.is(request.beforeSeq, -0))) {315
throw new RemoteError('gateway/bad-request', 'beforeSeq must be a non-negative safe integer', {})316
}317
validateHistoryWindow(request)318
}320
function validateHistoryWindow(request: Pick<SessionPageRequest, 'maxMessages' | 'turnWindow'>): void {321
if (request.maxMessages !== undefined322
&& (!Number.isSafeInteger(request.maxMessages) || request.maxMessages <= 0)) {323
throw new RemoteError('gateway/bad-request', 'maxMessages must be a positive safe integer', {})324
}325
const window = request.turnWindow326
if (window !== undefined) {327
if (!Number.isSafeInteger(window.minMessages) || window.minMessages <= 0328
|| window.minMessages > (request.maxMessages ?? DEFAULT_MAX_MESSAGES)) {329
throw new RemoteError('gateway/bad-request', 'turnWindow.minMessages must be a positive safe integer no greater than maxMessages', {})330
}331
if (!Number.isSafeInteger(window.minTurns) || window.minTurns <= 0) {332
throw new RemoteError('gateway/bad-request', 'turnWindow.minTurns must be a positive safe integer', {})333
}334
}335
}337
function addressId(address: SessionAddress): SessionId {338
return address.kind === 'session' ? address.sessionId : address.childSessionId339
}341
function validateAddress(342
address: SessionAddress,343
header: SessionHeader,344
inheritedEventCount: SessionLogOffsetType,345
projections: SessionObservation['projections'],346
): void {347
if (address.kind === 'session') {348
if (header.origin === 'subagent') {349
throw new RemoteError('session/agent-busy', 'subagent Sessions require their durable parent address', {350
reason: 'use subagent delivery for this child session',351
})352
}353
return354
}355
if (header.origin !== 'subagent' || header.parentSession !== address.parentSessionId) {356
throw new RemoteError('subagent/unauthorized', 'subagent does not belong to the supplied parent', {357
childSessionId: address.childSessionId,358
})359
}360
const identity = projections?.values.subagent361
if (identity === null) {362
throw new RemoteError('subagent/catalog-diagnostic', 'subagent descriptor is corrupt', {363
parentSessionId: address.parentSessionId,364
childSessionId: address.childSessionId,365
reason: 'corrupt',366
})367
}368
if (identity === undefined || identity.seq < inheritedEventCount) {369
throw new RemoteError('subagent/catalog-diagnostic', 'subagent descriptor is unavailable', {370
parentSessionId: address.parentSessionId,371
childSessionId: address.childSessionId,372
reason: 'unsupported',373
})374
}375
if (address.mode !== 'unknown' && identity.mode !== address.mode) {376
throw new RemoteError('subagent/unauthorized', 'subagent mode does not match the supplied address', {377
childSessionId: address.childSessionId,378
})379
}380
}382
function rejectNotFound(address: SessionAddress): never {383
if (address.kind === 'session') {384
throw new RemoteError('session/not-found', `session "${address.sessionId}" not found`, { sessionId: address.sessionId })385
}386
throw new RemoteError('subagent/not-found', 'subagent is unavailable', {387
parentSessionId: address.parentSessionId,388
childSessionId: address.childSessionId,389
})390
}392
function paginate(393
events: readonly SessionEvent[],394
beforeSeq: SessionLogOffsetType | undefined,395
maxMessages: number,396
throughSeq: SessionSeqCursor,397
turnWindow?: SessionPageRequest['turnWindow'],398
): { readonly events: SessionEvent[]; readonly hasMore: boolean } {399
const end = SessionLogOffset(Math.min(throughSeq + 1, beforeSeq ?? throughSeq + 1))400
let count = 0401
let turns = 0402
let cut = SessionLogOffset(0)403
for (let index = end - 1; index >= 0; index--) {404
const event = events[index] as SessionEvent405
if (turnWindow !== undefined && event.type === 'turn/start') {406
turns++407
if (count >= turnWindow.minMessages && turns >= turnWindow.minTurns) {408
cut = SessionLogOffset(index)409
break410
}411
}412
if (!MESSAGE_TYPES.has(event.type) || !isAppendSurfaceEvent(event)) continue413
count++414
const sources = event.sourceEventSeqs415
let groupStart = event.seq416
if (sources !== undefined) {417
for (const source of sources) {418
if (source < groupStart) groupStart = source419
}420
}421
if (count >= maxMessages) {422
cut = SessionLogOffset(groupStart)423
break424
}425
}426
return { events: events.slice(cut, end), hasMore: cut > 0 }427
}429
/** Translate current logical Session metadata to the browser wire. */430
function wireHeader(header: SessionHeader): SessionWireHeader {431
return { ...header }432
}434
function entryFor(event: SessionEvent): SessionEventEntry {435
return {436
type: 'event',437
// Session.append validates and freezes event data as JSON before publication.438
event: event as unknown as SessionWireEvent,439
}440
}442
/** Encode one bounded logical page without changing its pagination cut. */443
function pageRecords(events: readonly SessionEvent[]): SessionHistoryRecord[] {444
return events.map(entryFor)445
}