| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427 |
- /** Cold Session history pagination and live-event source. */
- import type { Context } from '@deepseek-ai/cordis'
- import { Deque } from '@deepseek-ai/dsh-deque'
- import type { AssistantStreamFrame } from '@deepseek-ai/dsh-agent'
- import {
- isAppendSurfaceEvent,
- SessionLogOffset,
- SessionSeq,
- } from '@deepseek-ai/dsh-session'
- import type {
- SessionEvent,
- SessionHeader,
- SessionId,
- SessionLogOffset as SessionLogOffsetType,
- SessionSeqCursor,
- } from '@deepseek-ai/dsh-session'
- import { SessionQueryError, type SessionObservation } from '@deepseek-ai/dsh-session-query'
- import type {} from '@deepseek-ai/dsh-subagent'
- import { RemoteError } from '@deepseek-ai/dsh-typert-protocol'
- import type { JsonValue } from '@deepseek-ai/dsh-util-values'
- import type {
- SessionAddress,
- SessionAssistantStreamFrame,
- SessionEventEntry,
- SessionFollowRequest,
- SessionFollowFrame,
- SessionHistoryRecord,
- SessionPage,
- SessionPageRequest,
- SessionProjectionBaseline,
- SessionProjectionValues,
- SessionWireHeader,
- SessionWireEvent,
- } from './types.ts'
- import { SessionAssistantStreamAccumulator } from './assistant-stream.ts'
- const DEFAULT_MAX_MESSAGES = 50
- const MESSAGE_TYPES = new Set(['user/message', 'assistant/message'])
- /** Implements cold-safe history operations delegated by the Session Controller. */
- export class SessionHistoryController {
- private readonly closeFollowers = new Set<() => void>()
- private readonly assistantStreams = new Map<SessionId, SessionAssistantStreamAccumulator>()
- /**
- * @param ctx - Host context carrying Session query and projection services.
- * @param promote - starts ordinary Session activation after snapshot delivery.
- */
- constructor(
- private readonly ctx: Context,
- private readonly promote: (observation: SessionObservation) => void,
- ) {
- ctx.on('agent/assistant-stream', ({ agent, frame }) => {
- let stream = this.assistantStreams.get(agent.session.id)
- if (stream === undefined) {
- stream = new SessionAssistantStreamAccumulator()
- this.assistantStreams.set(agent.session.id, stream)
- }
- stream.accept(frame, cursorBeforeNext(agent.session.seq))
- }, { global: true })
- ctx.on('agent/disposed', ({ agent }) => {
- this.assistantStreams.delete(agent.session.id)
- }, { global: true })
- ctx.effect(() => () => {
- for (const close of this.closeFollowers) close()
- this.closeFollowers.clear()
- }, 'session-controller.history')
- }
- /**
- * Read one message-aligned history page without activating an Agent.
- * @param request - durable address and backwards-page cursor.
- * @param signal - caller cancellation for persistence reads.
- * @returns a contiguous event page.
- */
- async page(request: SessionPageRequest, signal: AbortSignal): Promise<SessionPage> {
- validatePageRequest(request)
- const throughSeq: SessionSeqCursor = request.throughSeq === -1
- ? -1
- : SessionSeq(request.throughSeq)
- const beforeSeq = request.beforeSeq === undefined
- ? undefined
- : SessionLogOffset(request.beforeSeq)
- using source = await this.sourceFor(request.address, signal, false)
- signal.throwIfAborted()
- const sourceLog = source.events
- const sourceCursor: SessionSeqCursor = sourceLog.at(-1)?.seq ?? -1
- if (throughSeq > sourceCursor) {
- throw new RemoteError(
- 'gateway/bad-request',
- `session page through seq ${String(throughSeq)} is past cursor ${String(sourceCursor)}`,
- {},
- )
- }
- /* v8 ignore next -- Session and persistence validation guarantee a dense zero-based event prefix. */
- if (throughSeq >= 0 && sourceLog[throughSeq]?.seq !== throughSeq) {
- throw new RemoteError('gateway/internal', `session log does not contain through seq ${String(throughSeq)}`, {})
- }
- const page = paginate(
- sourceLog,
- beforeSeq,
- request.maxMessages ?? DEFAULT_MAX_MESSAGES,
- throughSeq,
- )
- const records = pageRecords(page.events)
- return {
- records,
- hasMore: page.hasMore,
- }
- }
- /**
- * Follow events appended after an initial cursor on one durable address.
- * @param request - durable address and last committed sequence already held by the caller.
- * @param signal - stream cancellation owned by the Remote carrier.
- * @returns a complete opening snapshot followed by gap-free durable events and opted-in assistant frames.
- */
- async *follow(request: SessionFollowRequest, signal: AbortSignal): AsyncIterable<SessionFollowFrame> {
- validateFollowRequest(request)
- const { address } = request
- const target = addressId(address)
- const buffered = new Deque<
- | { readonly type: 'event'; readonly event: SessionEvent }
- | {
- readonly type: 'assistant-stream'
- readonly frame: SessionAssistantStreamFrame
- readonly ordinal: number
- }
- >()
- let snapshotCursor: SessionSeqCursor | undefined
- let assistantStreamOrdinal = 0
- let wake: (() => void) | undefined
- const notify = (): void => {
- const resume = wake
- wake = undefined
- resume?.()
- }
- const follower = { closed: false }
- const close = (): void => {
- follower.closed = true
- notify()
- }
- this.closeFollowers.add(close)
- const disposeEvent = this.ctx.on('session/event', (session, event) => {
- if (session.id !== target) return
- buffered.pushBack({ type: 'event', event })
- notify()
- }, { global: true })
- const disposeCreated = this.ctx.on('session/created', (session) => {
- if (session.id !== target) return
- // Constructor seed events have no session/event notification. Normally
- // only the end-seed suffix is new; if persistence advanced after the
- // opening observation, replay everything beyond that snapshot cursor.
- const suffix = session.snapshotEvents(snapshotCursor === undefined
- ? session.firstLiveSeq
- : SessionLogOffset(snapshotCursor + 1))
- for (let index = suffix.length - 1; index >= 0; index -= 1) {
- buffered.pushFront({ type: 'event', event: suffix[index] as SessionEvent })
- }
- notify()
- }, { global: true })
- const disposeAssistantStream = request.assistantStream !== true
- ? undefined
- : this.ctx.on('agent/assistant-stream', ({ agent, frame }) => {
- if (agent.session.id !== target) return
- buffered.pushBack({
- type: 'assistant-stream',
- frame: wireAssistantStreamFrame(frame, cursorBeforeNext(agent.session.seq)),
- ordinal: ++assistantStreamOrdinal,
- })
- notify()
- }, { global: true })
- const onAbort = (): void => { notify() }
- signal.addEventListener('abort', onAbort, { once: true })
- try {
- using source = await this.sourceFor(address, signal, true)
- const events = source.events
- signal.throwIfAborted()
- const cursor = source.cursor
- snapshotCursor = cursor
- const page = paginate(events, undefined, request.maxMessages ?? DEFAULT_MAX_MESSAGES)
- const assistantStream = request.assistantStream === true
- ? this.assistantStreams.get(target)?.snapshot() ?? { revision: 0 }
- : undefined
- // The accumulator snapshot and this watermark are synchronous. Frames
- // through the cut are represented or superseded by that baseline,
- // including larger revisions from a retired Agent; later revision
- // resets reach Client continuity validation.
- const assistantStreamOrdinalCut = assistantStreamOrdinal
- yield {
- type: 'snapshot',
- header: wireHeader(source.header),
- cursor,
- records: pageRecords(page.events),
- hasMore: page.hasMore,
- projections: source.projections === undefined
- ? { asOfSeq: cursor, values: {} }
- : projectionBlock(source.projections),
- ...assistantStream === undefined ? {} : { assistantStream },
- }
- if (address.kind === 'session' && source.source === 'prepared') {
- const promotion = source.retain()
- try {
- this.promote(promotion)
- } catch (error: unknown) {
- promotion[Symbol.dispose]()
- throw error
- }
- }
- let nextOffset = SessionLogOffset(cursor + 1)
- while (!follower.closed && !signal.aborted) {
- const item = buffered.popFront()
- if (item === undefined) {
- await new Promise<void>((resolve) => { wake = resolve })
- continue
- }
- if (item.type === 'assistant-stream') {
- if (item.ordinal > assistantStreamOrdinalCut) {
- yield { type: 'assistant-stream', frame: item.frame }
- }
- continue
- }
- const expectedSeq = SessionSeq(nextOffset)
- if (item.event.seq < expectedSeq) continue
- if (item.event.seq !== expectedSeq) {
- throw new RemoteError('gateway/internal', `session event stream skipped seq ${String(expectedSeq)}`, {})
- }
- nextOffset = SessionLogOffset(nextOffset + 1)
- yield entryFor(item.event)
- }
- } finally {
- this.closeFollowers.delete(close)
- signal.removeEventListener('abort', onAbort)
- disposeCreated()
- disposeEvent()
- disposeAssistantStream?.()
- }
- }
- private async sourceFor(
- address: SessionAddress,
- signal: AbortSignal,
- withProjections: boolean,
- ): Promise<SessionObservation> {
- const sessionId = addressId(address)
- try {
- const observation = await this.ctx.sessionQuery.observeSession(sessionId, {
- signal,
- projectionMode: withProjections || address.kind === 'subagent' ? 'all' : 'none',
- })
- if (observation.header.cwd === undefined) {
- observation[Symbol.dispose]()
- rejectNotFound(address)
- }
- try {
- validateAddress(
- address,
- observation.header,
- observation.inheritedEventCount,
- observation.projections,
- )
- } catch (error: unknown) {
- observation[Symbol.dispose]()
- throw error
- }
- return observation
- } catch (error: unknown) {
- if (error instanceof SessionQueryError
- && error.code === 'SESSION_QUERY_SESSION_NOT_FOUND') rejectNotFound(address)
- throw error
- }
- }
- }
- function cursorBeforeNext(nextSeq: SessionLogOffsetType): SessionSeqCursor {
- return nextSeq === 0 ? -1 : SessionSeq(nextSeq - 1)
- }
- function wireAssistantStreamFrame(
- frame: AssistantStreamFrame,
- durableCursor: SessionSeqCursor,
- ): SessionAssistantStreamFrame {
- if (frame.type === 'start') return { ...frame, startedAfterSeq: durableCursor }
- if (frame.type === 'end') return frame
- return {
- ...frame,
- chunk: frame.chunk as JsonValue,
- }
- }
- function projectionBlock(
- snapshot: NonNullable<SessionObservation['projections']>,
- ): SessionProjectionBaseline {
- return {
- asOfSeq: snapshot.asOfSeq,
- // Projection definitions validate whole JSON values before snapshot publication.
- values: snapshot.values as SessionProjectionValues,
- }
- }
- function validatePageRequest(request: SessionPageRequest): void {
- if (!Number.isSafeInteger(request.throughSeq)
- || request.throughSeq < -1
- || Object.is(request.throughSeq, -0)) {
- throw new RemoteError('gateway/bad-request', 'throughSeq must be an integer greater than or equal to -1', {})
- }
- if (request.beforeSeq !== undefined
- && (!Number.isSafeInteger(request.beforeSeq)
- || request.beforeSeq < 0
- || Object.is(request.beforeSeq, -0))) {
- throw new RemoteError('gateway/bad-request', 'beforeSeq must be a non-negative safe integer', {})
- }
- if (request.maxMessages !== undefined
- && (!Number.isSafeInteger(request.maxMessages) || request.maxMessages <= 0)) {
- throw new RemoteError('gateway/bad-request', 'maxMessages must be a positive safe integer', {})
- }
- }
- function validateFollowRequest(request: SessionFollowRequest): void {
- if (request.maxMessages !== undefined
- && (!Number.isSafeInteger(request.maxMessages) || request.maxMessages <= 0)) {
- throw new RemoteError('gateway/bad-request', 'maxMessages must be a positive safe integer', {})
- }
- }
- function addressId(address: SessionAddress): SessionId {
- return address.kind === 'session' ? address.sessionId : address.childSessionId
- }
- function validateAddress(
- address: SessionAddress,
- header: SessionHeader,
- inheritedEventCount: SessionLogOffsetType,
- projections: SessionObservation['projections'],
- ): void {
- if (address.kind === 'session') {
- if (header.origin === 'subagent') {
- throw new RemoteError('session/agent-busy', 'subagent Sessions require their durable parent address', {
- reason: 'use subagent delivery for this child session',
- })
- }
- return
- }
- if (header.origin !== 'subagent' || header.parentSession !== address.parentSessionId) {
- throw new RemoteError('subagent/unauthorized', 'subagent does not belong to the supplied parent', {
- childSessionId: address.childSessionId,
- })
- }
- const identity = projections?.values.subagent
- if (identity === null) {
- throw new RemoteError('subagent/catalog-diagnostic', 'subagent descriptor is corrupt', {
- parentSessionId: address.parentSessionId,
- childSessionId: address.childSessionId,
- reason: 'corrupt',
- })
- }
- if (identity === undefined || identity.seq < inheritedEventCount) {
- throw new RemoteError('subagent/catalog-diagnostic', 'subagent descriptor is unavailable', {
- parentSessionId: address.parentSessionId,
- childSessionId: address.childSessionId,
- reason: 'unsupported',
- })
- }
- if (identity.mode !== address.mode) {
- throw new RemoteError('subagent/unauthorized', 'subagent mode does not match the supplied address', {
- childSessionId: address.childSessionId,
- })
- }
- }
- function rejectNotFound(address: SessionAddress): never {
- if (address.kind === 'session') {
- throw new RemoteError('session/not-found', `session "${address.sessionId}" not found`, { sessionId: address.sessionId })
- }
- throw new RemoteError('subagent/not-found', 'subagent is unavailable', {
- parentSessionId: address.parentSessionId,
- childSessionId: address.childSessionId,
- })
- }
- function paginate(
- events: readonly SessionEvent[],
- beforeSeq: SessionLogOffsetType | undefined,
- maxMessages: number,
- throughSeq: SessionSeqCursor = events.at(-1)?.seq ?? -1,
- ): { readonly events: SessionEvent[]; readonly hasMore: boolean } {
- const end = SessionLogOffset(Math.min(throughSeq + 1, beforeSeq ?? throughSeq + 1))
- let count = 0
- let cut = SessionLogOffset(0)
- for (let index = end - 1; index >= 0; index--) {
- const event = events[index] as SessionEvent
- if (!MESSAGE_TYPES.has(event.type) || !isAppendSurfaceEvent(event)) continue
- count++
- const sources = event.sourceEventSeqs
- let groupStart = event.seq
- if (sources !== undefined) {
- for (const source of sources) {
- if (source < groupStart) groupStart = source
- }
- }
- if (count >= maxMessages) {
- cut = SessionLogOffset(groupStart)
- break
- }
- }
- return { events: events.slice(cut, end), hasMore: cut > 0 }
- }
- /** Translate current logical Session metadata to the browser wire. */
- function wireHeader(header: SessionHeader): SessionWireHeader {
- return { ...header }
- }
- function entryFor(event: SessionEvent): SessionEventEntry {
- return {
- type: 'event',
- // Session.append validates and freezes event data as JSON before publication.
- event: event as unknown as SessionWireEvent,
- }
- }
- /** Encode one bounded logical page without changing its pagination cut. */
- function pageRecords(events: readonly SessionEvent[]): SessionHistoryRecord[] {
- return events.map(entryFor)
- }
|