| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248 |
- /**
- * Event-sourced session service: append-only session log, in-memory store, and
- * the derived LLM message history. Persistence is a plugin concern (subscribe
- * to `session/event`, drain on `session/flush`).
- *
- * @module @deepseek-ai/dsh-session
- */
- import { Context, Service } from '@deepseek-ai/cordis'
- import { isAbsolute } from 'node:path'
- import { brandString } from '@deepseek-ai/dsh-brand'
- import { assertNever, deepFreeze, snapshotJsonValue } from '@deepseek-ai/dsh-util-values'
- import { scopeOf, scopeTarget } from '@deepseek-ai/dsh-scope'
- import type { Scoped } from '@deepseek-ai/dsh-scope'
- import type { Message } from '@deepseek-ai/dsh-llm'
- import { SESSION_FORMAT_VERSION, SessionLogOffset, SessionSeq } from './types.ts'
- import type { TypertLookup } from '@deepseek-ai/dsh-typert-protocol'
- import type { CreateSessionOptions, EpochHeader, PrepareSessionOptions, RequestContext, SessionEvent, SessionEventMap, SessionEventType, SessionHeader, SessionId, SessionSeedEventState, SurfaceIntent, SurfaceEventType } from './types.ts'
- import { deriveEventMessage, SurfaceManager } from './surface.ts'
- import type { SessionSurface } from './surface.ts'
- import { foldRequestHeader } from './request-header.ts'
- export * from './types.ts'
- export { SessionPreparation } from './preparation.ts'
- export type { SessionPreparationOptions } from './preparation.ts'
- export type { AssistantMessage, ToolResultMessage, UserMessage } from '@deepseek-ai/dsh-llm'
- export { interruptedTurnClosers, TOOL_NOT_STARTED, TOOL_OUTCOME_UNKNOWN } from './repair.ts'
- export type { SessionSurface, SurfaceFoldReplacement, SurfaceFoldResult } from './surface.ts'
- export { deriveEventMessage, foldSurface, isAppendSurfaceEvent, isReplacementSurfaceEvent, isSurfaceEvent, isSurfaceEligibleType } from './surface.ts'
- export { canonicalHeader, foldRequestHeader, headerEquals } from './request-header.ts'
- export { KNOWN_SESSION_EVENT_TYPES } from './known-event-types.ts'
- declare module '@deepseek-ai/cordis' {
- interface Context {
- sessions: SessionStore
- }
- interface Events {
- /**
- * Creation announcement during session publication. A synchronous throw vetoes and rolls
- * back with a paired disposal; detach requested during dispatch is deferred.
- * A returned-promise rejection is logged but cannot retroactively veto this
- * synchronous boundary.
- * Scope-filtered dispatch (`@deepseek-ai/dsh-scope`): agent-scoped listeners
- * receive only sessions entered through that agent's context.
- * @param session - the session just entered and announced.
- * @dshScopeScan unsupported
- * @mode emit
- */
- 'session/created'(this: Scoped<Session>, session: Session): void
- /**
- * Emitted once when an announced session leaves the store, including
- * publication rollback, but never for an entry whose creation announcement
- * did not begin. Listener failures are logged and contained.
- * Scope-filtered dispatch (`@deepseek-ai/dsh-scope`) reuses the owner scope.
- * @param session - the session that is no longer live in the store.
- * @dshScopeScan unsupported
- * @mode emit
- */
- 'session/disposed'(this: Scoped<Session>, session: Session): void
- /**
- * Post-commit, fire-and-forget append feed. The listener snapshot resolves
- * before the log push, but callbacks run after it; observer failures are
- * logged and contained without making the committed append fail.
- * Scope-filtered dispatch (`@deepseek-ai/dsh-scope`): agent-scoped listeners
- * receive only events from sessions entered through that agent's context.
- * @param session - the session whose log grew.
- * @param event - the appended event, exactly as recorded.
- * @dshScopeScan unsupported
- * @mode emit
- */
- 'session/event'(this: Scoped<Session>, session: Session, event: SessionEvent): void
- /**
- * Awaited parallel durability checkpoint: every listener runs and the
- * caller awaits all of them, with no waterfall veto. Scope-filtered dispatch
- * (`@deepseek-ai/dsh-scope`) reuses the session's owner scope.
- * @param session - the session whose buffered events must reach durable storage.
- * @dshScopeScan unsupported
- * @mode parallel
- */
- 'session/flush'(this: Scoped<Session>, session: Session): Promise<void> | void
- }
- }
- declare module '@deepseek-ai/dsh-typert-protocol' {
- interface TypertLookupMap {
- session: TypertLookup<Session, SessionId>
- }
- }
- /** Validate and freeze one detached creation header in place. */
- function validateSessionHeader(id: SessionId, input: unknown): SessionHeader {
- if (input === null || typeof input !== 'object' || Array.isArray(input)) {
- throw new Error('session header is not a plain JSON record')
- }
- const record = input as Record<string, unknown>
- if (Object.hasOwn(record, 'seedLength')) {
- throw new Error('session header has invalid field "seedLength"')
- }
- if (record.version !== SESSION_FORMAT_VERSION) {
- throw new Error(`session header version must be ${SESSION_FORMAT_VERSION}, got ${String(record.version)}`)
- }
- if (record.id !== id) {
- throw new Error(`session header id "${String(record.id)}" does not match session id "${id}"`)
- }
- if (typeof record.createdAt !== 'number'
- || !Number.isSafeInteger(record.createdAt)
- || record.createdAt < 0) {
- throw new Error('session header createdAt must be a non-negative safe integer')
- }
- if (record.cwd !== undefined) {
- if (typeof record.cwd !== 'string') throw new Error('session header cwd must be a string')
- if (!isAbsolute(record.cwd)) {
- throw new Error(`session header cwd must be an absolute path, got "${record.cwd}"`)
- }
- }
- if (record.parentSession !== undefined && typeof record.parentSession !== 'string') {
- throw new Error('session header parentSession must be a string')
- }
- if (typeof record.isSeeded !== 'boolean') {
- throw new Error('session header isSeeded must be a boolean')
- }
- if (record.origin !== undefined && record.origin !== 'subagent') {
- throw new Error('session header origin must be "subagent"')
- }
- if (record.delegationDepth !== undefined
- && (typeof record.delegationDepth !== 'number' || !Number.isSafeInteger(record.delegationDepth) || record.delegationDepth < 0)) {
- throw new Error('session header delegationDepth must be a non-negative safe integer')
- }
- if (record.agentPreset !== undefined && typeof record.agentPreset !== 'string') {
- throw new Error('session header agentPreset must be a string')
- }
- return deepFreeze(record as unknown as SessionHeader)
- }
- /** Validate and freeze one exclusively owned persistence header in place. */
- function validateRestoredSessionHeader(id: SessionId, input: unknown): SessionHeader {
- if (input !== null && typeof input === 'object' && !Array.isArray(input)) {
- const prototype = Reflect.getPrototypeOf(input)
- if (prototype !== Object.prototype && prototype !== null) {
- throw new Error('session header is not a plain JSON record')
- }
- }
- return validateSessionHeader(id, input)
- }
- /** Detach, validate, and freeze the creation metadata published by a session. */
- function snapshotSessionHeader(id: SessionId, source?: SessionHeader): SessionHeader {
- const input: unknown = source === undefined
- ? { version: SESSION_FORMAT_VERSION, id, createdAt: Date.now(), isSeeded: false }
- : source
- const snapshot = snapshotJsonValue(input)
- if (snapshot === undefined) throw new Error('session header is not losslessly JSON-serializable')
- return validateSessionHeader(id, snapshot)
- }
- /**
- * Validate an exclusively owned event and deeply freeze its identified message
- * without copying the event. The caller transfers an object graph that no
- * producer retains and that shares no mutable children with another event.
- * Use {@link snapshotSessionEvent} when exclusive ownership is not guaranteed.
- * @param event - exclusively owned event imported across a trusted boundary.
- * @returns the same event object with a validated, deeply frozen message.
- */
- export function adoptSessionEvent<T extends SessionEvent>(event: T): T {
- assertMessageEventShape(
- event,
- `session event at seq ${event.seq}`,
- )
- switch (event.type) {
- case 'user/message':
- deepFreeze(event.data)
- break
- case 'assistant/message':
- case 'tool/result':
- deepFreeze(event.data.message)
- break
- default:
- // SessionEventMap is merge-extensible; plugin-owned events carry no core message.
- break
- }
- return event
- }
- /**
- * Detach one event while preserving deep immutability for its identified message.
- * @param event - event imported across a query or persistence boundary.
- * @returns a detached event snapshot with a validated, deeply frozen message.
- */
- export function snapshotSessionEvent<T extends SessionEvent>(event: T): T {
- return adoptSessionEvent(structuredClone(event))
- }
- /** Validate the fixed event envelope after one-pass JSON materialization. */
- function assertSessionEventEnvelope(value: Record<string, unknown>, index: number): asserts value is SessionEvent {
- const event = value
- for (const key in event) {
- switch (key) {
- case 'type':
- case 'seq':
- case 'time':
- case 'data':
- case 'surfaceOp':
- case 'sourceEventSeqs':
- case 'ignorable':
- break
- default:
- throw new Error(`seed event at index ${index} has an invalid event envelope`)
- }
- }
- const type = event['type']
- const seq = event['seq']
- const time = event['time']
- if (typeof type !== 'string'
- || typeof seq !== 'number' || !Number.isSafeInteger(seq) || seq < 0 || Object.is(seq, -0)
- || typeof time !== 'number' || !Number.isSafeInteger(time)
- || event['data'] === undefined
- || (event['ignorable'] !== undefined && event['ignorable'] !== true)) {
- throw new Error(`seed event at index ${index} has an invalid event envelope`)
- }
- switch (type) {
- case 'request/header':
- case 'user/message':
- case 'assistant/attempt':
- case 'assistant/message':
- case 'tool/result':
- assertCurrentLlmShape(event, index)
- break
- }
- }
- /** Reject obsolete request headers and malformed messages at the seed/load boundary. */
- function assertCurrentLlmShape(event: Record<string, unknown>, index: number): void {
- const data = event['data']
- const record = typeof data === 'object' && data !== null
- ? data as Record<string, unknown>
- : undefined
- if (event['type'] === 'request/header') {
- const header = record?.['header']
- const headerRecord = typeof header === 'object' && header !== null && !Array.isArray(header)
- ? header as Record<string, unknown>
- : undefined
- const config = headerRecord?.['config']
- if (!hasProviderModel(config)) throw new Error(`seed request/header at index ${index} lacks provider/model`)
- const configRecord = config as Record<string, unknown>
- const reasoningEffort = configRecord['reasoningEffort']
- if (reasoningEffort !== undefined
- && (typeof reasoningEffort !== 'string' || reasoningEffort.length === 0)) {
- throw new Error(`seed request/header at index ${index} has an invalid reasoningEffort`)
- }
- assertAdapterDefaults(headerRecord?.['adapterDefaults'], configRecord, index)
- const reason = record?.['reason']
- if (reason !== 'initial' && reason !== 'resume' && reason !== 'change' && reason !== 'series') {
- throw new Error(`seed request/header at index ${index} has an invalid reason`)
- }
- if (record?.['startsSeries'] !== undefined && record['startsSeries'] !== true) {
- throw new Error(`seed request/header at index ${index} has an invalid startsSeries marker`)
- }
- }
- const type = event['type']
- if (type === 'assistant/attempt') {
- assertAssistantSettlementShape(record, type, index)
- return
- }
- if (type !== 'user/message' && type !== 'assistant/message'
- && type !== 'tool/result') return
- assertMessageEventShape(event, `seed ${type} at index ${index}`)
- if (type === 'assistant/message') {
- assertAssistantSettlementShape(record, type, index)
- }
- }
- /** Validate fields used directly by restored Session lifecycle logic without replaying the embedded stream. */
- function assertAssistantSettlementShape(
- data: Record<string, unknown> | undefined,
- type: 'assistant/attempt' | 'assistant/message',
- index: number,
- ): void {
- const turn = data?.['turn']
- const step = data?.['step']
- if (typeof turn !== 'number' || !Number.isSafeInteger(turn) || turn < 0 || Object.is(turn, -0)
- || typeof step !== 'number' || !Number.isSafeInteger(step) || step < 0 || Object.is(step, -0)
- || !Array.isArray(data?.['stream'])) {
- throw new Error(`seed ${type} at index ${index} has invalid settlement fields`)
- }
- }
- const allowedAdapterKeys = new Set(['reasoningEffort', 'maxTokens'])
- /** Validate adapter-default markers imported from a durable request header. */
- function assertAdapterDefaults(
- value: unknown,
- config: Record<string, unknown>,
- index: number,
- ): void {
- if (value === undefined) return
- if (typeof value !== 'object' || value === null || Array.isArray(value)) {
- throw new Error(`seed request/header at index ${index} has invalid adapterDefaults`)
- }
- const defaults = value as Record<string, unknown>
- if (Object.keys(defaults).some(key => !allowedAdapterKeys.has(key))
- || Object.values(defaults).some(marker => marker !== true)
- || defaults['reasoningEffort'] === true && config['reasoningEffort'] === undefined
- || defaults['maxTokens'] === true && config['maxTokens'] === undefined) {
- throw new Error(`seed request/header at index ${index} has invalid adapterDefaults`)
- }
- }
- /** Validate only the event-specific invariants needed to safely replay a message. */
- function assertMessageEventShape(event: Record<string, unknown>, subject: string): void {
- const type = event['type']
- if (type !== 'user/message' && type !== 'assistant/message'
- && type !== 'tool/result') return
- const data = event['data']
- const record = typeof data === 'object' && data !== null
- ? data as Record<string, unknown>
- : undefined
- const message = type === 'user/message' ? record : record?.['message']
- if (typeof message !== 'object' || message === null
- || typeof (message as Record<string, unknown>)['id'] !== 'string'
- || (message as Record<string, unknown>)['id'] === '') {
- throw new Error(`${subject} lacks an identified message`)
- }
- const messageRecord = message as Record<string, unknown>
- const expectedRole = type === 'assistant/message' ? 'assistant' : 'user'
- if (messageRecord['role'] !== expectedRole) {
- throw new Error(`${subject} message must have role "${expectedRole}"`)
- }
- const source = messageRecord['source']
- if (typeof source !== 'object' || source === null
- || typeof (source as Record<string, unknown>)['kind'] !== 'string'
- || (source as Record<string, unknown>)['kind'] === '') {
- throw new Error(`${subject} message has invalid source`)
- }
- if (!Array.isArray(messageRecord['content'])) {
- throw new Error(`${subject} message has invalid content`)
- }
- const sourceRecord = source as Record<string, unknown>
- if (type === 'assistant/message') {
- if (sourceRecord['kind'] !== 'model' || !hasProviderModel(sourceRecord)) {
- throw new Error(`${subject} message must have model source`)
- }
- return
- }
- if (type !== 'tool/result') return
- if (sourceRecord['kind'] !== 'tool'
- || typeof sourceRecord['callId'] !== 'string'
- || sourceRecord['callId'] === '') {
- throw new Error(`${subject} message must have tool source`)
- }
- const content = messageRecord['content'] as unknown[]
- const block = content[0]
- if (content.length !== 1 || typeof block !== 'object' || block === null
- || (block as Record<string, unknown>)['type'] !== 'tool-result'
- || !Array.isArray((block as Record<string, unknown>)['content'])) {
- throw new Error(`${subject} message must contain one tool-result block`)
- }
- if ((block as Record<string, unknown>)['toolCallId'] !== sourceRecord['callId']) {
- throw new Error(`${subject} message has mismatched tool call ids`)
- }
- }
- /** Whether an unknown value carries the current provider/model pair. */
- function hasProviderModel(value: unknown): boolean {
- if (typeof value !== 'object' || value === null) return false
- const pair = value as Record<string, unknown>
- return typeof pair['provider'] === 'string' && pair['provider'].length > 0
- && typeof pair['model'] === 'string' && pair['model'].length > 0
- }
- type SessionCallback = (...args: unknown[]) => unknown
- /** Resolve one listener snapshot, including Cordis's internal dispatch checks. */
- function collectSessionCallbacks(ctx: Context, args: unknown[]): SessionCallback[] {
- return [...ctx.events.dispatch('emit', args)] as SessionCallback[]
- }
- /** Invoke one resolved observe-only listener snapshot with per-listener containment. */
- function invokeContainedSessionObservers(
- ctx: Context,
- name: 'session/event' | 'session/disposed',
- id: SessionId,
- args: unknown[],
- callbacks: SessionCallback[],
- ): void {
- for (const callback of callbacks) {
- try {
- const returned: unknown = callback(...args)
- void Promise.resolve(returned).catch((error: unknown) => {
- ctx.logger.warn(`session "${id}": ${name} listener rejected: ${String(error)}`)
- })
- } catch (error: unknown) {
- ctx.logger.warn(`session "${id}": ${name} listener threw: ${String(error)}`)
- }
- }
- }
- /** All mutable lifecycle state for one exact store entry. */
- interface SessionEntry {
- readonly id: SessionId
- readonly session: Session
- readonly carrier: Scoped<Session>
- readonly emitCtx: Context
- announced: boolean
- announcing: boolean
- appending: boolean
- detachRequested: boolean
- detach(): void
- }
- /** Store attachment for the append path; module-private to keep Session store-agnostic publicly. */
- const attachments = new WeakMap<Session, SessionEntry>()
- /**
- * An event-sourced session: an append-only log of {@link SessionEvent}s.
- *
- * Plain class (not a Service) — create live instances via
- * `ctx.sessions.create()` and detached instances via {@link create}.
- * Seeding with an existing event log replays/forks a session.
- * @typert object
- */
- export class Session {
- private log: SessionEvent[] = []
- /** Single incremental owner of surface acceptance and projection state. */
- private readonly surfaceManager = new SurfaceManager(this.log)
- /** The ordered surface over this session's event log. */
- get surface(): SessionSurface {
- return this.surfaceManager
- }
- /**
- * Detached, deep-frozen creation metadata (format version, cwd, lineage,
- * and whether fork history exists). Supplied by the store via `ctx.sessions.create()`. When a
- * `Session` is created without a store-owned header, a minimal header is
- * synthesized (stamped with the current {@link SESSION_FORMAT_VERSION}) so
- * `session.header` is always present. Kept out of the event log — it is a
- * storage concern, not replayable conversation state.
- */
- readonly header: SessionHeader
- /** Number of leading events inherited from this Session's fork parent. */
- readonly inheritedEventCount: SessionLogOffset
- /** The session identity, derived from its durable header's single copy. */
- get id(): SessionId {
- return this.header.id
- }
- /**
- * The first seq appended IN THIS PROCESS: the length of the constructor
- * seed (0 without one). Events with smaller seq values entered through
- * construction — replay, fork, or resume — and were never published on the
- * `session/event` firehose (constructor seeds do not emit). This offset marks
- * the constructor-input boundary for lifecycle ownership and persistence
- * adoption; consumers that need complete canonical history still start at
- * seq 0. Distinct from {@link inheritedEventCount}, the DURABLE
- * fork-lineage cut: a resumed session's constructor seed is its full stored
- * log, while the inherited count keeps the original fork value — this field is the
- * in-process construction fact.
- *
- * Not persisted itself: a seeded session projects it into the log as the
- * `session/end-seed` event, which is what a consumer reading STORED history
- * reads. Locate the LAST such event, not necessarily one at this seq — a
- * seed already ending in one is not re-marked, so reopening an untouched
- * session leaves that event at a smaller seq than `firstLiveSeq`. Prefer
- * this field in-process: it is exact before the marker reaches storage.
- *
- * When this lifecycle appends the marker, it occupies this seq before the
- * store attaches and therefore does not publish either. Otherwise this seq
- * holds an ordinary published write.
- */
- readonly firstLiveSeq: SessionLogOffset
- /**
- * Create a detached session by validating and snapshotting borrowed seed
- * events and storage metadata.
- * @param id - session identity.
- * @param seed - optional borrowed replay or fork events.
- * @param header - optional borrowed storage metadata.
- * @param inheritedEventCount - exact fork-inherited prefix length for a seeded header.
- * @returns a detached session.
- */
- static create(
- id: SessionId,
- seed?: readonly SessionEvent[],
- header?: SessionHeader,
- inheritedEventCount?: SessionLogOffset,
- ): Session {
- return new Session(id, seed, header, 'snapshot', inheritedEventCount)
- }
- /**
- * Restore a detached session by adopting an independently owned or deeply frozen seed.
- * Runtime-required event fields, event envelopes, sequence continuity, surface
- * transitions, and header fields are validated without copying or freezing events.
- * Embedded Assistant streams remain opaque until a stream consumer or storage
- * verifier reads them.
- * @param id - restored session identity.
- * @param seed - independently owned or deeply frozen events.
- * @param header - independently owned storage metadata.
- * @param inheritedEventCount - exact fork-inherited prefix length decoded from storage.
- * @param eventState - aliasing state carried from the operation that produced the seed.
- * @returns a restored detached session.
- */
- static fromRestore(
- id: SessionId,
- seed: readonly SessionEvent[],
- header: SessionHeader,
- inheritedEventCount: SessionLogOffset,
- eventState: SessionSeedEventState,
- ): Session {
- return new Session(
- id,
- seed,
- header,
- eventState,
- inheritedEventCount,
- )
- }
- private constructor(
- id: SessionId,
- seed?: readonly SessionEvent[],
- header?: SessionHeader,
- mode: 'snapshot' | SessionSeedEventState = 'snapshot',
- suppliedInheritedEventCount?: SessionLogOffset,
- ) {
- const restoredHeader = mode === 'snapshot' ? undefined : validateRestoredSessionHeader(id, header)
- if (seed !== undefined) {
- // Validate the seed to the SAME invariants `append` enforces, so a
- // replay/fork (`ctx.sessions.create(id, { seed })`) cannot construct a
- // live log that no persistence backend could store: each event's `data`
- // must be JSON-serializable, and `seq` must be contiguous from 0 (the
- // `seq = log.length` contract the whole system relies on). Without this,
- // a bad seed would surface only later as a backend rejection or a silent
- // divergence between the live log and disk.
- for (const [index, source] of seed.entries()) {
- // The seed is a persistence/replay boundary: validate and detach the
- // complete event in one lossless-JSON pass.
- const snapshot = mode === 'snapshot' ? snapshotJsonValue(source) : source
- if (snapshot === undefined) {
- throw new Error(`seed event at index ${index} is not losslessly JSON-serializable`)
- }
- assertSessionEventEnvelope(snapshot, index)
- if (snapshot.seq !== index) {
- throw new Error(`seed event at index ${index} has seq ${snapshot.seq} (expected ${index}); seed must be contiguous from 0`)
- }
- // A seed is accepted incrementally through the same transition as a
- // live append and a full-log fold. The candidate is planned before it
- // enters `log`, so a failure cannot partially mutate the surface.
- try {
- this.surfaceManager.validateNext(snapshot)
- } catch (error: unknown) {
- throw new Error(`invalid seed event at index ${index}: ${error instanceof Error ? error.message : 'invalid surface metadata'}`)
- }
- this.log.push(mode === 'snapshot' ? deepFreeze(snapshot) : snapshot)
- }
- }
- this.firstLiveSeq = SessionLogOffset(this.log.length)
- this.header = restoredHeader ?? snapshotSessionHeader(id, header)
- if (this.header.isSeeded && seed === undefined) {
- throw new Error('seeded session requires an explicit constructor seed')
- }
- if (this.header.isSeeded && suppliedInheritedEventCount === undefined) {
- throw new Error('seeded session requires an inherited event count')
- }
- const inheritedEventCount = SessionLogOffset(suppliedInheritedEventCount ?? 0)
- if (!this.header.isSeeded && inheritedEventCount !== 0) {
- throw new Error('unseeded session inherited event count must be 0')
- }
- if (inheritedEventCount > this.log.length) {
- throw new Error('session inherited event count exceeds its event log')
- }
- if (mode === 'snapshot' && this.header.isSeeded && inheritedEventCount !== this.log.length) {
- throw new Error('seeded session constructor seed must equal its inherited prefix')
- }
- this.inheritedEventCount = inheritedEventCount
- // A fresh seeded child always owns one tagged marker at its inherited cut,
- // even when the copied prefix already ends in an ancestor marker. Restore
- // retains that durable marker and appends only the ordinary resume marker.
- if (seed !== undefined && mode === 'snapshot' && this.header.isSeeded) {
- this.append('session/end-seed', { inherited: true })
- } else if (seed !== undefined && this.log.at(-1)?.type !== 'session/end-seed') {
- this.append('session/end-seed', {})
- }
- }
- /** Cached immutable full snapshot of the private append-only log. */
- private eventsSnapshot: readonly SessionEvent[] | undefined
- /**
- * Return the immutable event stored at one exact sequence number.
- * @param seq - event sequence number.
- * @returns the accepted event, or undefined when the log does not contain it.
- */
- eventAt(seq: SessionSeq): SessionEvent | undefined {
- return this.log[seq]
- }
- /**
- * Materialize an immutable snapshot of a half-open event sequence range.
- * A full current snapshot is reused until the next append; every previously
- * returned snapshot remains stable after later appends.
- * @param fromSeq - non-negative inclusive sequence number; defaults to the log start.
- * @param toSeqExclusive - non-negative exclusive sequence number; defaults to the current end.
- * @returns a frozen array of the selected deeply frozen events.
- */
- snapshotEvents(
- fromSeq: SessionLogOffset = SessionLogOffset(0),
- toSeqExclusive: SessionLogOffset = this.seq,
- ): readonly SessionEvent[] {
- if (fromSeq === 0 && toSeqExclusive === this.log.length) {
- this.eventsSnapshot ??= Object.freeze([...this.log])
- return this.eventsSnapshot
- }
- return Object.freeze(this.log.slice(fromSeq, toSeqExclusive))
- }
- /**
- * Return this Session's events after its fork-inherited prefix.
- * @returns a fresh array containing child-owned events in log order.
- */
- ownEvents(): readonly SessionEvent[] {
- return this.snapshotEvents(this.inheritedEventCount)
- }
- /**
- * Whether one existing event position is outside the fork-inherited prefix.
- * @param seq - event position in this Session.
- * @returns true when the event belongs to this Session rather than its parent.
- */
- isOwnSeq(seq: SessionSeq): boolean {
- return seq >= this.inheritedEventCount && seq < this.seq
- }
- /** The next event's sequence number — always the log length (the `seq = log.length` contiguity contract). */
- get seq(): SessionLogOffset {
- return SessionLogOffset(this.log.length)
- }
- /**
- * Append one typed event to the log and synchronously notify observers via
- * the store-owned, module-private publication hooks. The hot path never blocks
- * on I/O — persistence plugins buffer asynchronously. Once the event enters
- * the log, the append is committed: observer failures are logged and
- * contained per listener, so they do not change the return value or prevent
- * later listeners from observing the same accepted event.
- *
- * @param type - The event type (key of {@link SessionEventMap}).
- * @param data - The event payload; must be JSON-serializable.
- * @param opts - Surface metadata: `surfaceOp` controls how the event enters
- * the ordered surface; `sourceEventSeqs` lists the seq numbers of earlier
- * events this one derives from. REQUIRED for
- * {@link SurfaceEventType} events (every message-producing event must
- * declare how it joins the surface, the sole source of derived model
- * history) and
- * rejected by the compiler for non-surface types like `turn/start` or
- * `assistant/attempt`. Assistant messages embed their exact provider
- * stream and cannot cite top-level source events.
- * @returns the logged event — its assigned `seq`/`time` plus the SNAPSHOT of
- * `data` that entered the log, so reading `event.data` back sees the logged
- * value, never the caller's still-mutable input.
- * @throws if `data` or surface metadata is not losslessly JSON-serializable
- * (BigInt, function, symbol, undefined, negative zero, non-finite number,
- * circular reference, sparse array, or an exotic object such as
- * Map/Set/Date/class instance), or when the candidate violates the
- * canonical surface contract (marker shape and eligibility, unique
- * earlier source-event references, positional replacement validity, and complete
- * shadowed-node coverage). One iterative pass reads, validates, and
- * copies each nested value once, so a stateful getter cannot supply one value
- * to validation and another to storage. The event log is the durable source
- * of truth, so a bad event fails at the append site rather than later during
- * a backend flush. A synchronous internal dispatch validation failure or an
- * append reentered while this acceptance/publication boundary is open also
- * rejects before the log changes.
- */
- append<T extends SessionEventType>(
- type: T,
- data: SessionEventMap[T],
- ...opts: T extends SurfaceEventType ? [opts: SurfaceIntent<T>] : []
- ): SessionEvent<T> {
- const surfaceOpts: SurfaceIntent | undefined = opts[0]
- const surfaceMetadata = {
- ...surfaceOpts?.sourceEventSeqs === undefined ? {} : { sourceEventSeqs: surfaceOpts.sourceEventSeqs },
- ...surfaceOpts?.surfaceOp === undefined ? {} : { surfaceOp: surfaceOpts.surfaceOp },
- }
- const dataSnapshot = snapshotJsonValue(data)
- if (dataSnapshot === undefined) {
- throw new Error(`session event "${type}" carries non-JSON-serializable data`)
- }
- const surfaceMetadataSnapshot = snapshotJsonValue(surfaceMetadata)
- if (surfaceMetadataSnapshot === undefined) {
- throw new Error(`session event "${type}" carries non-JSON-serializable surface metadata`)
- }
- const entry = attachments.get(this)
- if (entry?.appending) {
- throw new Error('session append cannot reenter while another append is being published')
- }
- const event = deepFreeze({
- type,
- seq: SessionSeq(this.log.length),
- time: Date.now(),
- data: dataSnapshot,
- ...(surfaceMetadataSnapshot as { surfaceOp?: unknown; sourceEventSeqs?: unknown }),
- } as unknown as SessionEvent<T>)
- this.surfaceManager.validateNext(event as SessionEvent)
- if (entry !== undefined) entry.appending = true
- try {
- let callbacks: SessionCallback[] | undefined
- const callbackArgs: unknown[] = [this, event]
- if (entry !== undefined) {
- callbacks = collectSessionCallbacks(entry.emitCtx, [entry.carrier, 'session/event', ...callbackArgs])
- }
- this.log.push(event as SessionEvent)
- this.eventsSnapshot = undefined
- if (callbacks !== undefined && entry !== undefined) {
- invokeContainedSessionObservers(entry.emitCtx, 'session/event', entry.id, callbackArgs, callbacks)
- }
- return event
- } finally {
- if (entry !== undefined) {
- entry.appending = false
- if (entry.detachRequested && !entry.announcing) entry.detach()
- }
- }
- }
- /** Cached fold of the request-header events — see {@link requestHeader}. */
- private headerFold: EpochHeader | undefined
- /** Log position (events consumed) the header fold has reached. */
- private headerFoldSeq = 0
- /**
- * The {@link EpochHeader} in force after the log's last header event — the
- * header the NEXT request will be compared against — or undefined before
- * the first `request/header` snapshot. The live, incrementally-maintained
- * form of `foldRequestHeader(session.snapshotEvents())`: each header event is folded
- * once, when first seen, so a per-step read costs O(new events).
- * @returns the folded header, or undefined when no header event exists yet.
- */
- requestHeader(): EpochHeader | undefined {
- if (this.headerFoldSeq < this.log.length) {
- // Frozen on update: the fold is session state exposed by reference — a
- // consumer mutating it in place (instead of building a replacement)
- // would desync every later comparison against the log, so mutation
- // throws instead.
- this.headerFold = deepFreeze(foldRequestHeader(this.log.slice(this.headerFoldSeq), this.headerFold))
- this.headerFoldSeq = this.log.length
- }
- return this.headerFold
- }
- /** Cached fold of `request/context` events. */
- private contextFold: RequestContext | undefined
- private contextFoldSeq = 0
- /**
- * Return the latest resolved route metadata, or `undefined` before the first
- * `request/context` event. Each event is folded once.
- * @returns the latest immutable route metadata.
- */
- requestContext(): RequestContext | undefined {
- if (this.contextFoldSeq < this.log.length) {
- for (const event of this.log.slice(this.contextFoldSeq)) {
- if (event.type === 'request/context') this.contextFold = deepFreeze({ ...event.data })
- }
- this.contextFoldSeq = this.log.length
- }
- return this.contextFold
- }
- /** The derived-message cache: frozen projections, extended per unseen node. */
- private derived: Message[] = []
- /** Surface position (nodes projected) the cache has reached. */
- private derivedNodes = 0
- /** {@link SurfaceManager.replaceGeneration} the cache was built under. */
- private derivedGeneration = 0
- /**
- * Derive the LLM message history by walking the ordered sequences of
- * message-producing events maintained by `surfaceOp` markers. The
- * surface is the single source of derived history: every message-producing
- * append records its `surfaceOp`, so a raw event with no marker (a chunk, a
- * turn boundary) is correctly absent, and a compaction `replace` deletes the
- * shadowed nodes from the derivation. The projection rules are
- * {@link deriveEventMessage}, folded per node.
- *
- * CACHED: each surface node is projected exactly once, when first seen — a
- * call costs O(new nodes), and a surface rewrite (a `replace`;
- * {@link SessionSurface.replaceGeneration}) rebuilds. The returned array is
- * a fresh snapshot per call (later appends never grow an array a caller
- * already holds); the `Message` objects in it are SHARED and **deep-frozen**.
- * Their content reuses the already frozen durable event data, so the cache
- * needs no second deep clone and consumers still cannot mutate the log.
- * @returns a fresh array of the shared, frozen derived history.
- */
- deriveMessages(): Message[] {
- const surface = this.surface
- const nodes = surface.nodes
- const generation = surface.replaceGeneration
- if (generation !== this.derivedGeneration) {
- this.derived = []
- this.derivedNodes = 0
- this.derivedGeneration = generation
- }
- for (const seq of nodes.slice(this.derivedNodes)) {
- // Surface sequences are built from this.log — seq is always a valid
- // index by construction. The non-null assertion expresses that invariant.
- // oxlint-disable-next-line typescript/no-non-null-assertion
- const msg = this.deriveEventMessage(this.log[seq]!)
- // A surface node is one of the five message-producing types, but an
- // empty-content assistant/message (a max-tokens step that hosts only
- // usage) derives to null and must not enter the transcript.
- if (msg) this.derived.push(msg)
- }
- this.derivedNodes = nodes.length
- return [...this.derived]
- }
- /**
- * Instance face of the pure per-node `deriveEventMessage` export from
- * `surface.ts`.
- * @param event - the event to project.
- * @returns the derived message, or null when the event produces none.
- */
- deriveEventMessage(event: SessionEvent): Message | null {
- return deriveEventMessage(event)
- }
- }
- /** A fork source: either the live session object or its live store id. */
- export type SessionForkSource = Session | SessionId
- /**
- * Rejection codes for session forking: the fork source id is unknown to the
- * live store (`SESSION_NOT_FOUND`) or names a session object that is not the
- * store's live instance (`SESSION_NOT_LIVE`); the requested child id is
- * already taken (`SESSION_ALREADY_EXISTS`); the boundary is not a contiguous
- * existing seq (`INVALID_BOUNDARY`); or the selected prefix ends inside an
- * open turn (`OPEN_TURN`).
- */
- export type SessionForkErrorCode =
- | 'SESSION_NOT_FOUND'
- | 'SESSION_NOT_LIVE'
- | 'SESSION_ALREADY_EXISTS'
- | 'INVALID_BOUNDARY'
- | 'OPEN_TURN'
- /** Typed error for session fork rejections. */
- export class SessionForkError extends Error {
- constructor(message: string, public readonly code: SessionForkErrorCode) {
- super(message)
- this.name = 'SessionForkError'
- }
- }
- /**
- * In-memory session store (`ctx.sessions`).
- *
- * Persistence is intentionally not implemented here — the agent lifecycle
- * attaches a session-log writer to each published session's write handle;
- * a session published outside that lifecycle persists nothing.
- */
- export class SessionStore extends Service {
- private store = new Map<SessionId, SessionEntry>()
- private counter = 0
- constructor(ctx: Context) {
- super(ctx, 'sessions')
- ctx.inject(['typert'], (typeCtx) => {
- typeCtx.typert.lookups.register('session', {
- parameter: 'session',
- wire: 'sessionId',
- hostTypeSymbol: '@deepseek-ai/dsh-session#Session',
- wireTypeSymbol: '@deepseek-ai/dsh-session/types#SessionId',
- resolve: sessionId => this.get(sessionId),
- })
- })
- }
- /**
- * Create a session owned by the calling fiber: disposing that fiber stops
- * event notification and removes the session from the store. `options.seed`
- * populates the session with a copy of those events (replay/fork);
- * `options.meta` attaches creation metadata (validated absolute `cwd`, seed
- * and parent lineage, and delegation depth) as the immutable
- * {@link SessionHeader} (the store fills `version`/`id`/`createdAt`).
- *
- * For an agent whose session must be torn down IN ORDER with its loop (so the
- * loop's final events are published before the store attachment ends), do NOT use this
- * — fold the session lifecycle into the agent's own effect via
- * {@link prepare} + {@link enter} + {@link announce} (see
- * `dsh-agent-loop`'s creation transaction).
- *
- * @param id - the session id; omitted, the store mints `session-<n>`.
- * @param options - seed events and/or creation metadata for the header.
- * @returns the live session, already entered and announced.
- * @throws if a session with `id` already exists, metadata is not a plain
- * lossless-JSON record with valid scalar fields, or `meta.cwd` is a
- * non-absolute path (storage backends key directories off it).
- */
- create(id?: SessionId, options?: CreateSessionOptions): Session {
- const session = this.prepare(id, options)
- // Single effect owned by the calling fiber. Yield the detach BEFORE
- // announcing so a throwing `session/created` listener rolls the attach back
- // (the generator effect disposes already-yielded disposers on a throw)
- // instead of leaking the store entry and its publication hooks.
- this.ctx.effect(function* (this: SessionStore) {
- yield this.enter(session)
- this.announce(session)
- }.bind(this), 'sessions.create()')
- return session
- }
- /**
- * Build a session WITHOUT entering it into the store — validate the id/cwd and
- * construct the {@link Session} (with its immutable {@link SessionHeader}).
- * Pairs with {@link enter} + {@link announce}: a caller that owns a composite
- * `ctx.effect` (the agent factory) folds the session lifecycle into that ONE
- * effect so a fiber unload tears the session + agent down as a single ORDERED
- * chain rather than as racing sibling effects — which would remove the publication hooks
- * before the driver's closing events commit, dropping them.
- *
- * @param id - the session id; omitted, the store mints `session-<n>`.
- * @param options - seed events and/or creation metadata for the header. With
- * `eventState`, every seed event is either independently owned or any
- * shared value is deeply frozen; {@link Session.fromRestore} validates and
- * adopts those values without copying or freezing them.
- * @returns the constructed session, NOT yet in the store.
- * @throws if a session with `id` already exists, metadata is not a plain
- * lossless-JSON record with valid scalar fields, or `meta.cwd` is a
- * non-absolute path.
- */
- prepare(id?: SessionId, options?: PrepareSessionOptions): Session {
- let sessionId: SessionId
- if (id === undefined) {
- do sessionId = brandString<SessionId>(`session-${++this.counter}`)
- while (this.store.has(sessionId))
- } else {
- sessionId = brandString<SessionId>(id)
- }
- if (this.store.has(sessionId)) throw new Error(`session "${sessionId}" already exists`)
- if (options !== undefined) {
- const { eventState } = options
- switch (eventState) {
- case 'detached':
- case 'shared-frozen':
- return Session.fromRestore(
- sessionId,
- options.seed,
- options.meta,
- options.inheritedEventCount,
- eventState,
- )
- case undefined:
- break
- /* v8 ignore next -- closed-union exhaustiveness guard */
- default:
- assertNever(eventState, 'SessionStore.prepare event state')
- }
- }
- const seed = options?.seed
- const meta = options?.meta
- const header: SessionHeader = {
- version: SESSION_FORMAT_VERSION,
- id: sessionId,
- createdAt: meta?.createdAt ?? Date.now(),
- ...meta?.cwd === undefined ? {} : { cwd: meta.cwd },
- ...meta?.parentSession === undefined ? {} : { parentSession: meta.parentSession },
- isSeeded: meta?.isSeeded ?? false,
- ...meta?.origin === undefined ? {} : { origin: meta.origin },
- ...meta?.delegationDepth === undefined ? {} : { delegationDepth: meta.delegationDepth },
- ...meta?.agentPreset === undefined ? {} : { agentPreset: meta.agentPreset },
- }
- return Session.create(sessionId, seed, header, options?.inheritedEventCount)
- }
- /**
- * Enter a {@link prepare}d session into the store: install the module-private
- * append publication hooks and add it to the store. Returns the DETACH
- * disposer (hooks + store removal). Does NOT emit `session/created` —
- * the caller yields this disposer inside its effect and THEN calls
- * {@link announce}, so a throwing `session/created` listener rolls the attach
- * back instead of leaking it.
- *
- * Re-checks the id for a duplicate: `prepare` and `enter` are public
- * cross-package primitives and a caller may interleave arbitrary work (or
- * another create) between them, so a stale prepared session must NOT overwrite
- * a live store entry of the same id — its detach disposer would later delete
- * the REAL session. The {@link create} convenience and the agent factory call
- * the two back-to-back so they never trip this, but the public API cannot
- * assume that.
- *
- * @param session - a {@link prepare}d session not yet in the store.
- * @returns the detach disposer (publication hooks + store removal). When called from
- * a synchronous `session/created` listener, removal and disposal wait until
- * that creation dispatch unwinds.
- * @throws if a session with this id is already in the store.
- */
- enter(session: Session): () => void {
- const id = session.id
- const carrier = scopeTarget(session, scopeOf(this.ctx))
- // This is the authoritative collision boundary after arbitrary unpublished
- // preparation. Only one exact same-id transaction can publish.
- if (this.store.has(id)) throw new Error(`session "${id}" already exists`)
- if (attachments.has(session)) throw new Error(`session "${id}" is already attached to a store`)
- const entry: SessionEntry = {
- id,
- session,
- carrier,
- emitCtx: this.ctx,
- announced: false,
- announcing: false,
- appending: false,
- detachRequested: false,
- detach: () => { this.detachEntered(entry) },
- }
- this.store.set(id, entry)
- attachments.set(session, entry)
- let entered = true
- const detach = (): void => {
- if (!entered) return
- entered = false
- // A lifecycle listener may own the advanced detach capability. Keep the
- // entry and its publication hooks live until synchronous creation or append
- // publication unwinds, then publish the paired disposal edge.
- if (entry.announcing || entry.appending) {
- entry.detachRequested = true
- return
- }
- entry.detach()
- }
- return detach
- }
- /** Remove one exact entered session and emit its paired disposal when announced. */
- private detachEntered(entry: SessionEntry): void {
- entry.detachRequested = false
- // A stale capability cannot remove observers or storage belonging to a
- // later same-id lifecycle.
- /* v8 ignore next -- enter() rejects replacement while this single-shot detach capability is live. */
- if (this.store.get(entry.id) !== entry) return
- this.store.delete(entry.id)
- attachments.delete(entry.session)
- if (entry.announced) this.emitDisposed(entry)
- }
- /** Emit `session/created` exactly once for an {@link enter}ed session (with
- * the carrier {@link enter} captured). Separate from {@link enter} so the
- * caller can yield the detach disposer first (rollback safety — see
- * {@link enter}).
- * @param session - the entered session to announce to listeners.
- * @throws if the session is not live or its announcement already began,
- * including a reentrant call from a creation listener. */
- announce(session: Session): void {
- const entry = this.liveEntryFor(session)
- if (entry.announced || entry.announcing) {
- throw new Error(`session "${entry.id}" was already announced`)
- }
- // Mark before emit: Cordis emit may deliver to earlier listeners and then
- // throw. Rollback must still pair that partial creation with disposal, and
- // a listener cannot recursively create a second lifecycle edge.
- entry.announced = true
- const callbackArgs: unknown[] = [session]
- entry.announcing = true
- try {
- const callbacks = collectSessionCallbacks(this.ctx, [entry.carrier, 'session/created', session])
- for (const callback of callbacks) {
- // Synchronous throws intentionally propagate and veto publication; the
- // yielded detach then emits the paired disposal edge. An async function
- // is nevertheless assignable to a void listener, so observe its returned
- // promise: rejection is too late to roll back and must be logged instead
- // of becoming unhandled.
- const returned: unknown = callback(...callbackArgs)
- void Promise.resolve(returned).catch((error: unknown) => {
- this.ctx.logger.warn(`session "${entry.id}": session/created listener rejected: ${String(error)}`)
- })
- }
- } finally {
- entry.announcing = false
- if (entry.detachRequested && !entry.appending) entry.detach()
- }
- }
- /** Emit the paired teardown notification with per-listener containment. */
- private emitDisposed(entry: SessionEntry): void {
- const callbackArgs: unknown[] = [entry.session]
- try {
- const callbacks = collectSessionCallbacks(this.ctx, [entry.carrier, 'session/disposed', entry.session])
- invokeContainedSessionObservers(this.ctx, 'session/disposed', entry.id, callbackArgs, callbacks)
- } catch (error: unknown) {
- this.ctx.logger.warn(`session "${entry.id}": session/disposed dispatch threw: ${String(error)}`)
- }
- }
- /**
- * Dispatch the awaited `session/flush` durability checkpoint for `session`,
- * with the carrier captured at {@link enter}. THE flush entry point: the
- * store owns the carrier, so callers (the checkpoint policy's per-request
- * barrier, goal-round-driver's idle checkpoint, teardown drains, and consumers
- * that flush themselves before reading storage) must come through here
- * rather than dispatch a raw `ctx.parallel('session/flush', …)` — one owner,
- * one spelling, and the scoped-dispatch invariant can pin it.
- * @param session - the session whose buffered events must reach durable storage.
- * @returns whether at least one durability listener participated, after every
- * listener has settled successfully.
- * @throws the first registered listener failure after every listener settles.
- */
- async flush(session: Session): Promise<boolean> {
- const { carrier } = this.liveEntryFor(session)
- const callbackArgs: unknown[] = [session]
- const callbacks = collectSessionCallbacks(this.ctx, [carrier, 'session/flush', session])
- const results = await Promise.allSettled(callbacks.map((callback) => {
- try {
- return callback(...callbackArgs)
- } catch (error: unknown) {
- // Preserve the listener's exact rejection value; flush is a caller-owned
- // failure boundary, and Cordis listeners may throw arbitrary values.
- // oxlint-disable-next-line typescript/prefer-promise-reject-errors
- return Promise.reject(error)
- }
- }))
- const failure = results.find((result): result is PromiseRejectedResult => result.status === 'rejected')
- if (failure !== undefined) throw failure.reason
- return callbacks.length > 0
- }
- /** Return the exact live entry; detached/prepared objects reject. */
- private liveEntryFor(session: Session): SessionEntry {
- const entry = attachments.get(session)
- if (entry === undefined || this.store.get(entry.id) !== entry) {
- throw new Error(`session "${session.id}" is not live in this store`)
- }
- return entry
- }
- /**
- * Look up a live session.
- * @param id - the session id to look up.
- * @returns the session, or undefined when no live session has that id.
- */
- get(id: SessionId): Session | undefined {
- return this.store.get(id)?.session
- }
- /**
- * All live sessions, in creation order.
- * @returns a fresh array; mutating it does not affect the store.
- */
- list(): Session[] {
- return [...this.store.values()].map(entry => entry.session)
- }
- /**
- * Create a live child session from a stable prefix of a live source.
- * `boundary` is an inclusive source event seq; omitted means the source's
- * current last event. The selected slice may end with a between-turn event
- * but must not end inside an open turn.
- *
- * @param source - Live source session object or id.
- * @param boundary - Inclusive source event seq to fork through; omitted means
- * the source's current last event, and omitted on an empty source forks an
- * empty child.
- * @param childSessionId - Optional child session id; omitted delegates to
- * `SessionStore`'s id policy.
- * @returns The created live child session.
- */
- fork(source: SessionForkSource, boundary?: SessionSeq, childSessionId?: SessionId): Session {
- if (childSessionId !== undefined && this.get(childSessionId) !== undefined) {
- throw new SessionForkError(`session "${childSessionId}" already exists`, 'SESSION_ALREADY_EXISTS')
- }
- const liveSource = this._resolveForkSource(source)
- const seed = this._forkSeed(liveSource, boundary)
- return this.create(childSessionId, {
- seed,
- inheritedEventCount: SessionLogOffset(seed.length),
- meta: {
- ...liveSource.header.cwd !== undefined ? { cwd: liveSource.header.cwd } : {},
- parentSession: liveSource.id,
- isSeeded: true,
- },
- })
- }
- private _forkSeed(session: Session, requestedBoundary: SessionSeq | undefined): readonly SessionEvent[] {
- const lastEvent = session.snapshotEvents().at(-1)
- let boundary: SessionSeq
- if (requestedBoundary !== undefined) {
- boundary = requestedBoundary
- } else {
- if (lastEvent === undefined) return []
- boundary = lastEvent.seq
- }
- if (!Number.isSafeInteger(boundary) || boundary < 0) {
- throw new SessionForkError(
- `fork boundary for session "${session.id}" must be a non-negative safe integer, got ${String(boundary)}`,
- 'INVALID_BOUNDARY',
- )
- }
- if (boundary >= session.seq) {
- const lastSeq = lastEvent?.seq
- throw new SessionForkError(
- `fork boundary ${boundary} does not exist in session "${session.id}" (last seq: ${lastSeq ?? 'none'})`,
- 'INVALID_BOUNDARY',
- )
- }
- const boundaryEvent = session.eventAt(boundary)
- if (boundaryEvent === undefined || boundaryEvent.seq !== boundary) {
- throw new SessionForkError(
- `fork boundary ${boundary} does not match a contiguous event seq in session "${session.id}"`,
- 'INVALID_BOUNDARY',
- )
- }
- const events = session.snapshotEvents(SessionLogOffset(0), SessionLogOffset(boundary + 1))
- const lastTurnBoundary = events
- .findLast(event => event.type === 'turn/start' || event.type === 'turn/end')
- if (lastTurnBoundary?.type === 'turn/start') {
- throw new SessionForkError(
- `fork boundary ${boundary} in session "${session.id}" ends inside open turn ${lastTurnBoundary.data.turn}`,
- 'OPEN_TURN',
- )
- }
- return events
- }
- private _resolveForkSource(source: SessionForkSource): Session {
- if (typeof source === 'string') {
- const session = this.get(source)
- if (session === undefined) throw new SessionForkError(`session "${source}" not found`, 'SESSION_NOT_FOUND')
- return session
- }
- const live = this.get(source.id)
- if (live === undefined) {
- throw new SessionForkError(`session "${source.id}" not found`, 'SESSION_NOT_FOUND')
- }
- if (live !== source) throw new SessionForkError(`session "${source.id}" is not the live store instance`, 'SESSION_NOT_LIVE')
- return source
- }
- }
- export { decodeSeqRanges, encodeSeqRanges } from './seq-ranges.ts'
- export default SessionStore
|