返回源码地图

packages/core/session/src/index.ts

main snapshot · da00f7f5358f · 正文引用章节 03 / 05 / 05 / 05 / 05 / 05 / 05 / 05 / 05 / 05 / 06;完整原文可核对,不声称全文件人工逐行审计

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

1/**
2 * Event-sourced session service: append-only session log, in-memory store, and
3 * the derived LLM message history. Persistence is a plugin concern (subscribe
4 * to `session/event`, drain on `session/flush`).
5 *
6 * @module @deepseek-ai/dsh-session
7 */
8
9import { Context, Service } from '@deepseek-ai/cordis'
10import { isAbsolute } from 'node:path'
11import { brandString } from '@deepseek-ai/dsh-brand'
12import { assertNever, deepFreeze, snapshotJsonValue } from '@deepseek-ai/dsh-util-values'
13import { scopeOf, scopeTarget } from '@deepseek-ai/dsh-scope'
14import type { Scoped } from '@deepseek-ai/dsh-scope'
15import type { Message } from '@deepseek-ai/dsh-llm'
16import { SESSION_FORMAT_VERSION, SessionLogOffset, SessionSeq } from './types.ts'
17import type { TypertLookup } from '@deepseek-ai/dsh-typert-protocol'
18import type { CreateSessionOptions, EpochHeader, PrepareSessionOptions, RequestContext, SessionEvent, SessionEventMap, SessionEventType, SessionHeader, SessionId, SessionSeedEventState, SurfaceIntent, SurfaceEventType } from './types.ts'
19import { SurfaceManager, validateSessionEventData, validateSurfaceMetadata } from './surface.ts'
20import type { SessionSurface, SessionMessageProjection } from './surface.ts'
21import { foldRequestHeader } from './request-header.ts'
22import { ToolHistoryProjection } from './tool-history.ts'
23import type { ToolHistory } from '@deepseek-ai/dsh-llm'
24
25import { buildForkSeed } from './fork.ts'
26
27export { buildForkSeed } from './fork.ts'
28export * from './types.ts'
29export { SessionPreparation } from './preparation.ts'
30export type { SessionPreparationOptions } from './preparation.ts'
31export type { AssistantMessage, DeveloperMessage, SystemMessage, ToolResultMessage, UserMessage } from '@deepseek-ai/dsh-llm'
32export { interruptedTurnClosers, ToolCallRecovery, TOOL_NOT_STARTED, TOOL_OUTCOME_UNKNOWN } from './repair.ts'
33export type { SessionSurface, SurfaceFoldReplacement, SurfaceFoldResult, SessionMessageProjection, SessionMessageProjectionContext } from './surface.ts'
34export { deriveEventMessage, foldSurface, isAppendSurfaceEvent, isReplacementSurfaceEvent, isSurfaceEvent, isSurfaceEligibleType } from './surface.ts'
35export { canonicalHeader, foldRequestHeader, headerEquals } from './request-header.ts'
36export { KNOWN_SESSION_EVENT_TYPES } from './known-event-types.ts'
37
38declare module '@deepseek-ai/cordis' {
39 interface Context {
40 sessions: SessionStore
41 }
42
43 interface Events {
44 /**
45 * Creation announcement during session publication. A synchronous throw vetoes and rolls
46 * back with a paired disposal; detach requested during dispatch is deferred.
47 * A returned-promise rejection is logged but cannot retroactively veto this
48 * synchronous boundary.
49 * Scope-filtered dispatch (`@deepseek-ai/dsh-scope`): agent-scoped listeners
50 * receive only sessions entered through that agent's context.
51 * @param session - the session just entered and announced.
52 * @mode emit
53 */
54 'session/created'(this: Scoped<Session>, session: Session): void
55 /**
56 * Emitted once when an announced session leaves the store, including
57 * publication rollback, but never for an entry whose creation announcement
58 * did not begin. Listener failures are logged and contained.
59 * Scope-filtered dispatch (`@deepseek-ai/dsh-scope`) reuses the owner scope.
60 * @param session - the session that is no longer live in the store.
61 * @mode emit
62 */
63 'session/disposed'(this: Scoped<Session>, session: Session): void
64 /**
65 * Post-commit, fire-and-forget append feed. The listener snapshot resolves
66 * before the log push, but callbacks run after it; observer failures are
67 * logged and contained without making the committed append fail.
68 * Scope-filtered dispatch (`@deepseek-ai/dsh-scope`): agent-scoped listeners
69 * receive only events from sessions entered through that agent's context.
70 * @param session - the session whose log grew.
71 * @param event - the appended event, exactly as recorded.
72 * @mode emit
73 */
74 'session/event'(this: Scoped<Session>, session: Session, event: SessionEvent): void
75 /**
76 * Awaited parallel durability checkpoint: every listener runs and the
77 * caller awaits all of them, with no waterfall veto. Scope-filtered dispatch
78 * (`@deepseek-ai/dsh-scope`) reuses the session's owner scope.
79 * @param session - the session whose buffered events must reach durable storage.
80 * @mode parallel
81 */
82 'session/flush'(this: Scoped<Session>, session: Session): Promise<void> | void
83 }
84}
85
86declare module '@deepseek-ai/dsh-typert-protocol' {
87 interface TypertLookupMap {
88 session: TypertLookup<Session, SessionId>
89 }
90}
91
92/** Validate and freeze one detached creation header in place. */
93function validateSessionHeader(id: SessionId, input: unknown): SessionHeader {
94 if (input === null || typeof input !== 'object' || Array.isArray(input)) {
95 throw new Error('session header is not a plain JSON record')
96 }
97 const record = input as Record<string, unknown>
98 if (Object.hasOwn(record, 'seedLength')) {
99 throw new Error('session header has invalid field "seedLength"')
100 }
101 if (record.version !== SESSION_FORMAT_VERSION) {
102 throw new Error(`session header version must be ${SESSION_FORMAT_VERSION}, got ${String(record.version)}`)
103 }
104 if (record.id !== id) {
105 throw new Error(`session header id "${String(record.id)}" does not match session id "${id}"`)
106 }
107 if (typeof record.createdAt !== 'number'
108 || !Number.isSafeInteger(record.createdAt)
109 || record.createdAt < 0) {
110 throw new Error('session header createdAt must be a non-negative safe integer')
111 }
112 if (record.cwd !== undefined) {
113 if (typeof record.cwd !== 'string') throw new Error('session header cwd must be a string')
114 if (!isAbsolute(record.cwd)) {
115 throw new Error(`session header cwd must be an absolute path, got "${record.cwd}"`)
116 }
117 }
118 if (record.parentSession !== undefined && typeof record.parentSession !== 'string') {
119 throw new Error('session header parentSession must be a string')
120 }
121 if (typeof record.isSeeded !== 'boolean') {
122 throw new Error('session header isSeeded must be a boolean')
123 }
124 if (record.origin !== undefined && record.origin !== 'subagent') {
125 throw new Error('session header origin must be "subagent"')
126 }
127 if (record.delegationDepth !== undefined
128 && (typeof record.delegationDepth !== 'number' || !Number.isSafeInteger(record.delegationDepth) || record.delegationDepth < 0)) {
129 throw new Error('session header delegationDepth must be a non-negative safe integer')
130 }
131 if (record.agentPreset !== undefined && typeof record.agentPreset !== 'string') {
132 throw new Error('session header agentPreset must be a string')
133 }
134 return deepFreeze(record as unknown as SessionHeader)
135}
136
137/** Validate and freeze one exclusively owned persistence header in place. */
138function validateRestoredSessionHeader(id: SessionId, input: unknown): SessionHeader {
139 if (input !== null && typeof input === 'object' && !Array.isArray(input)) {
140 const prototype = Reflect.getPrototypeOf(input)
141 if (prototype !== Object.prototype && prototype !== null) {
142 throw new Error('session header is not a plain JSON record')
143 }
144 }
145 return validateSessionHeader(id, input)
146}
147
148/** Detach, validate, and freeze the creation metadata published by a session. */
149function snapshotSessionHeader(id: SessionId, source?: SessionHeader): SessionHeader {
150 const input: unknown = source === undefined
151 ? { version: SESSION_FORMAT_VERSION, id, createdAt: Date.now(), isSeeded: false }
152 : source
153 const snapshot = snapshotJsonValue(input)
154 if (snapshot === undefined) throw new Error('session header is not losslessly JSON-serializable')
155 return validateSessionHeader(id, snapshot)
156}
157
158/**
159 * Validate an exclusively owned event and deeply freeze its identified message
160 * without copying the event. The caller transfers an object graph that no
161 * producer retains and that shares no mutable children with another event.
162 * Use {@link snapshotSessionEvent} when exclusive ownership is not guaranteed.
163 * @param event - exclusively owned event imported across a trusted boundary.
164 * @returns the same event object with a validated, deeply frozen message.
165 * @throws when event-local surface metadata, request-header fields, or message invariants are invalid; history relations are not checked.
166 */
167export function adoptSessionEvent<T extends SessionEvent>(event: T): T {
168 validateSessionEventData(event, `session event at seq ${event.seq}`)
169 validateSurfaceMetadata(event)
170 assertMessageEventShape(
171 event,
172 `session event at seq ${event.seq}`,
173 )
174 switch (event.type) {
175 case 'user/message':
176 deepFreeze(event.data)
177 break
178 case 'developer/message':
179 case 'system/message':
180 case 'assistant/message':
181 case 'tool/result':
182 deepFreeze(event.data.message)
183 break
184 default:
185 // SessionEventMap is merge-extensible; plugin-owned events carry no core message.
186 break
187 }
188 return event
189}
190
191/**
192 * Detach one event while preserving deep immutability for its identified message.
193 * @param event - event imported across a query or persistence boundary.
194 * @returns a detached event snapshot with a validated, deeply frozen message.
195 */
196export function snapshotSessionEvent<T extends SessionEvent>(event: T): T {
197 return adoptSessionEvent(structuredClone(event))
198}
199
200/** Validate the fixed event envelope after one-pass JSON materialization. */
201function assertSessionEventEnvelope(value: unknown, index: number): asserts value is SessionEvent {
202 if (value === null || typeof value !== 'object' || Array.isArray(value)) {
203 throw new Error(`seed event at index ${index} has an invalid event envelope`)
204 }
205 const event = value as Record<string, unknown>
206 for (const key in event) {
207 switch (key) {
208 case 'type':
209 case 'seq':
210 case 'time':
211 case 'data':
212 case 'surfaceOp':
213 case 'sourceEventSeqs':
214 case 'ignorable':
215 break
216 default:
217 throw new Error(`seed event at index ${index} has an invalid event envelope`)
218 }
219 }
220 const type = event['type']
221 const seq = event['seq']
222 const time = event['time']
223 if (typeof type !== 'string'
224 || typeof seq !== 'number' || !Number.isSafeInteger(seq) || seq < 0 || Object.is(seq, -0)
225 || typeof time !== 'number' || !Number.isSafeInteger(time)
226 || event['data'] === undefined
227 || (event['ignorable'] !== undefined && event['ignorable'] !== true)) {
228 throw new Error(`seed event at index ${index} has an invalid event envelope`)
229 }
230 validateSessionEventData(event as SessionEvent, `seed ${type} at index ${index}`)
231 switch (type) {
232 case 'request/header':
233 case 'developer/message':
234 case 'system/message':
235 case 'user/message':
236 case 'assistant/attempt':
237 case 'assistant/message':
238 case 'tool/result':
239 assertCurrentLlmShape(event, index)
240 break
241 }
242}
243
244/** Reject obsolete request headers and malformed messages at the seed/load boundary. */
245function assertCurrentLlmShape(event: Record<string, unknown>, index: number): void {
246 const data = event['data']
247 const record = typeof data === 'object' && data !== null
248 ? data as Record<string, unknown>
249 : undefined
250 if (event['type'] === 'request/header') {
251 const headerRecord = record?.['header'] as Record<string, unknown>
252 const config = headerRecord['config']
253 if (!hasProviderModel(config)) throw new Error(`seed request/header at index ${index} lacks provider/model`)
254 const configRecord = config as Record<string, unknown>
255 const reasoningEffort = configRecord['reasoningEffort']
256 if (reasoningEffort !== undefined
257 && (typeof reasoningEffort !== 'string' || reasoningEffort.length === 0)) {
258 throw new Error(`seed request/header at index ${index} has an invalid reasoningEffort`)
259 }
260 assertAdapterDefaults(headerRecord['adapterDefaults'], configRecord, index)
261 const reason = record?.['reason']
262 if (reason !== 'initial' && reason !== 'resume' && reason !== 'change' && reason !== 'series') {
263 throw new Error(`seed request/header at index ${index} has an invalid reason`)
264 }
265 if (record?.['startsSeries'] !== undefined && record['startsSeries'] !== true) {
266 throw new Error(`seed request/header at index ${index} has an invalid startsSeries marker`)
267 }
268 }
269 const type = event['type']
270 if (type === 'assistant/attempt') {
271 assertAssistantSettlementShape(record, type, index)
272 return
273 }
274 if (!isMessageEventType(type)) return
275 assertMessageEventShape(event, `seed ${type} at index ${index}`)
276 if (type === 'assistant/message') {
277 assertAssistantSettlementShape(record, type, index)
278 }
279}
280
281/** Validate fields used directly by restored Session lifecycle logic without replaying the embedded stream. */
282function assertAssistantSettlementShape(
283 data: Record<string, unknown> | undefined,
284 type: 'assistant/attempt' | 'assistant/message',
285 index: number,
286): void {
287 const turn = data?.['turn']
288 const step = data?.['step']
289 if (typeof turn !== 'number' || !Number.isSafeInteger(turn) || turn < 0 || Object.is(turn, -0)
290 || typeof step !== 'number' || !Number.isSafeInteger(step) || step < 0 || Object.is(step, -0)
291 || !Array.isArray(data?.['stream'])) {
292 throw new Error(`seed ${type} at index ${index} has invalid settlement fields`)
293 }
294}
295
296const allowedAdapterKeys = new Set(['reasoningEffort', 'maxTokens'])
297
298/** Validate adapter-default markers imported from a durable request header. */
299function assertAdapterDefaults(
300 value: unknown,
301 config: Record<string, unknown>,
302 index: number,
303): void {
304 if (value === undefined) return
305 if (typeof value !== 'object' || value === null || Array.isArray(value)) {
306 throw new Error(`seed request/header at index ${index} has invalid adapterDefaults`)
307 }
308 const defaults = value as Record<string, unknown>
309 if (Object.keys(defaults).some(key => !allowedAdapterKeys.has(key))
310 || Object.values(defaults).some(marker => marker !== true)
311 || defaults['reasoningEffort'] === true && config['reasoningEffort'] === undefined
312 || defaults['maxTokens'] === true && config['maxTokens'] === undefined) {
313 throw new Error(`seed request/header at index ${index} has invalid adapterDefaults`)
314 }
315}
316
317/** The surface event types whose payload carries an identified message. */
318function isMessageEventType(type: unknown): type is SurfaceEventType {
319 return type === 'developer/message' || type === 'system/message' || type === 'user/message'
320 || type === 'assistant/message' || type === 'tool/result'
321}
322
323const MESSAGE_ROLE_BY_TYPE: Record<SurfaceEventType, Message['role']> = {
324 'system/message': 'system',
325 'developer/message': 'developer',
326 'user/message': 'user',
327 'assistant/message': 'assistant',
328 'tool/result': 'tool',
329}
330
331/** Validate only the event-specific invariants needed to safely replay a message. */
332function assertMessageEventShape(event: Record<string, unknown>, subject: string): void {
333 const type = event['type']
334 if (!isMessageEventType(type)) return
335 const data = event['data']
336 const record = typeof data === 'object' && data !== null
337 ? data as Record<string, unknown>
338 : undefined
339 const message = type === 'user/message' ? record : record?.['message']
340 if (typeof message !== 'object' || message === null
341 || typeof (message as Record<string, unknown>)['id'] !== 'string'
342 || (message as Record<string, unknown>)['id'] === '') {
343 throw new Error(`${subject} lacks an identified message`)
344 }
345 const messageRecord = message as Record<string, unknown>
346 const expectedRole = MESSAGE_ROLE_BY_TYPE[type]
347 if (messageRecord['role'] !== expectedRole) {
348 throw new Error(`${subject} message must have role "${expectedRole}"`)
349 }
350 const source = messageRecord['source']
351 if (typeof source !== 'object' || source === null
352 || typeof (source as Record<string, unknown>)['kind'] !== 'string'
353 || (source as Record<string, unknown>)['kind'] === '') {
354 throw new Error(`${subject} message has invalid source`)
355 }
356 if (!Array.isArray(messageRecord['content'])) {
357 throw new Error(`${subject} message has invalid content`)
358 }
359 const sourceRecord = source as Record<string, unknown>
360 if (type === 'system/message') {
361 if (sourceRecord['kind'] !== 'system-prompt') {
362 throw new Error(`${subject} message must have system-prompt source`)
363 }
364 return
365 }
366 if (type === 'assistant/message') {
367 if (sourceRecord['kind'] !== 'model' || !hasProviderModel(sourceRecord)) {
368 throw new Error(`${subject} message must have model source`)
369 }
370 return
371 }
372 if (type !== 'tool/result') return
373 if (sourceRecord['kind'] !== 'tool'
374 || typeof sourceRecord['callId'] !== 'string'
375 || sourceRecord['callId'] === '') {
376 throw new Error(`${subject} message must have tool source`)
377 }
378 if (messageRecord['toolCallId'] !== sourceRecord['callId']) {
379 throw new Error(`${subject} message has mismatched tool call ids`)
380 }
381}
382
383/** Whether an unknown value carries the current provider/model pair. */
384function hasProviderModel(value: unknown): boolean {
385 if (typeof value !== 'object' || value === null) return false
386 const pair = value as Record<string, unknown>
387 return typeof pair['provider'] === 'string' && pair['provider'].length > 0
388 && typeof pair['model'] === 'string' && pair['model'].length > 0
389}
390
391type SessionCallback = (...args: unknown[]) => unknown
392
393/** Resolve one listener snapshot, including Cordis's internal dispatch checks. */
394function collectSessionCallbacks(ctx: Context, args: unknown[]): SessionCallback[] {
395 return [...ctx.events.dispatch('emit', args)] as SessionCallback[]
396}
397
398/** Invoke one resolved observe-only listener snapshot with per-listener containment. */
399function invokeContainedSessionObservers(
400 ctx: Context,
401 name: 'session/event' | 'session/disposed',
402 id: SessionId,
403 args: unknown[],
404 callbacks: SessionCallback[],
405): void {
406 for (const callback of callbacks) {
407 try {
408 const returned: unknown = callback(...args)
409 void Promise.resolve(returned).catch((error: unknown) => {
410 ctx.logger.warn(`session "${id}": ${name} listener rejected: ${String(error)}`)
411 })
412 } catch (error: unknown) {
413 ctx.logger.warn(`session "${id}": ${name} listener threw: ${String(error)}`)
414 }
415 }
416}
417
418/** All mutable lifecycle state for one exact store entry. */
419interface SessionEntry {
420 readonly id: SessionId
421 readonly session: Session
422 readonly carrier: Scoped<Session>
423 readonly emitCtx: Context
424 announced: boolean
425 announcing: boolean
426 appending: boolean
427 detachRequested: boolean
428 detach(): void
429}
430
431/** Store attachment for the append path; module-private to keep Session store-agnostic publicly. */
432const attachments = new WeakMap<Session, SessionEntry>()
433
434/**
435 * An event-sourced session: an append-only log of {@link SessionEvent}s.
436 *
437 * Plain class (not a Service) — create live instances via
438 * `ctx.sessions.create()` and detached instances via {@link create}.
439 * Seeding with an existing event log replays/forks a session.
440 * @typert object
441 */
442export class Session {
443 private log: SessionEvent[] = []
444 /** Single incremental owner of surface acceptance and projection state. */
445 private readonly surfaceManager: SurfaceManager
446
447 /** The ordered surface over this session's event log. */
448 get surface(): SessionSurface {
449 return this.surfaceManager
450 }
451
452 /**
453 * Detached, deep-frozen creation metadata (format version, cwd, lineage,
454 * and whether fork history exists). Supplied by the store via `ctx.sessions.create()`. When a
455 * `Session` is created without a store-owned header, a minimal header is
456 * synthesized (stamped with the current {@link SESSION_FORMAT_VERSION}) so
457 * `session.header` is always present. Kept out of the event log — it is a
458 * storage concern, not replayable conversation state.
459 */
460 readonly header: SessionHeader
461
462 /** Number of leading events inherited from this Session's fork parent. */
463 readonly inheritedEventCount: SessionLogOffset
464
465 /** The session identity, derived from its durable header's single copy. */
466 get id(): SessionId {
467 return this.header.id
468 }
469
470 /**
471 * The constructor seed length (0 without one), before any marker appended
472 * during construction. Seed events never publish on `session/event`. A
473 * marker appended before the store attaches occupies this seq without
474 * publishing either; otherwise this seq is available for the next append.
475 *
476 * This in-process offset is not persisted. A fork seed can already contain
477 * the child's inherited marker and synthetic closers, so its child-owned
478 * history starts at {@link inheritedEventCount}, before this offset. A
479 * resumed Session's seed contains its full stored log, while its inherited
480 * count keeps the durable fork cut. Consumers needing complete canonical
481 * history start at seq 0.
482 */
483 readonly firstLiveSeq: SessionLogOffset
484
485 /**
486 * First event produced for this object lifecycle. A new fork includes its
487 * child-owned seed marker and closers; a restored Session starts after its
488 * complete stored prefix. This in-process capture offset is not persisted.
489 */
490 readonly firstLifecycleSeq: SessionLogOffset
491
492 /**
493 * Create a detached session by validating and snapshotting borrowed seed
494 * events and storage metadata.
495 * @param id - session identity.
496 * @param seed - optional borrowed replay or fork events.
497 * @param header - optional borrowed storage metadata.
498 * @param inheritedEventCount - exact fork-inherited prefix length for a seeded header.
499 * @param projections - pure interpreters for plugin-owned message changes.
500 * @returns a detached session.
501 * @throws when a seed event requires a missing message interpreter or fails validation.
502 */
503 static create(
504 id: SessionId,
505 seed?: readonly SessionEvent[],
506 header?: SessionHeader,
507 inheritedEventCount?: SessionLogOffset,
508 projections?: readonly SessionMessageProjection[],
509 ): Session {
510 return new Session(id, seed, header, 'snapshot', inheritedEventCount, projections)
511 }
512
513 /**
514 * Restore a detached session by adopting an independently owned or deeply frozen seed.
515 * Runtime-required event fields, event envelopes, sequence continuity, surface
516 * transitions, and header fields are validated without copying or freezing events.
517 * Embedded Assistant streams remain opaque until a stream consumer or storage
518 * verifier reads them.
519 * @param id - restored session identity.
520 * @param seed - independently owned or deeply frozen events.
521 * @param header - independently owned storage metadata.
522 * @param inheritedEventCount - exact fork-inherited prefix length decoded from storage.
523 * @param eventState - aliasing state carried from the operation that produced the seed.
524 * @param projections - pure interpreters for plugin-owned message changes.
525 * @returns a restored detached session.
526 * @throws when a seed event requires a missing message interpreter or fails validation.
527 */
528 static fromRestore(
529 id: SessionId,
530 seed: readonly SessionEvent[],
531 header: SessionHeader,
532 inheritedEventCount: SessionLogOffset,
533 eventState: SessionSeedEventState,
534 projections?: readonly SessionMessageProjection[],
535 ): Session {
536 return new Session(
537 id,
538 seed,
539 header,
540 eventState,
541 inheritedEventCount,
542 projections,
543 )
544 }
545
546 private constructor(
547 id: SessionId,
548 seed?: readonly SessionEvent[],
549 header?: SessionHeader,
550 mode: 'snapshot' | SessionSeedEventState = 'snapshot',
551 suppliedInheritedEventCount?: SessionLogOffset,
552 projections: readonly SessionMessageProjection[] = [],
553 ) {
554 this.surfaceManager = new SurfaceManager(this.log, SessionLogOffset(0), projections)
555 const restoredHeader = mode === 'snapshot' ? undefined : validateRestoredSessionHeader(id, header)
556 if (seed !== undefined) {
557 // Validate the seed to the SAME invariants `append` enforces, so a
558 // replay/fork (`ctx.sessions.create(id, { seed })`) cannot construct a
559 // live log that no persistence backend could store: each event's `data`
560 // must be JSON-serializable, and `seq` must be contiguous from 0 (the
561 // `seq = log.length` contract the whole system relies on). Without this,
562 // a bad seed would surface only later as a backend rejection or a silent
563 // divergence between the live log and disk.
564 for (const [index, source] of seed.entries()) {
565 // The seed is a persistence/replay boundary: validate and detach the
566 // complete event in one lossless-JSON pass.
567 const snapshot = mode === 'snapshot' ? snapshotJsonValue(source) : source
568 if (snapshot === undefined) {
569 throw new Error(`seed event at index ${index} is not losslessly JSON-serializable`)
570 }
571 assertSessionEventEnvelope(snapshot, index)
572 if (snapshot.seq !== index) {
573 throw new Error(`seed event at index ${index} has seq ${snapshot.seq} (expected ${index}); seed must be contiguous from 0`)
574 }
575 // A seed is accepted incrementally through the same transition as a
576 // live append and a full-log fold. The candidate is planned before it
577 // enters `log`, so a failure cannot partially mutate the surface.
578 try {
579 this.surfaceManager.validateNext(snapshot)
580 } catch (error: unknown) {
581 throw new Error(`invalid seed event at index ${index}: ${error instanceof Error ? error.message : 'invalid surface metadata'}`)
582 }
583 this.log.push(mode === 'snapshot' ? deepFreeze(snapshot) : snapshot)
584 }
585 }
586 this.firstLiveSeq = SessionLogOffset(this.log.length)
587 this.header = restoredHeader ?? snapshotSessionHeader(id, header)
588 if (this.header.isSeeded && seed === undefined) {
589 throw new Error('seeded session requires an explicit constructor seed')
590 }
591 if (this.header.isSeeded && suppliedInheritedEventCount === undefined) {
592 throw new Error('seeded session requires an inherited event count')
593 }
594 const inheritedEventCount = SessionLogOffset(suppliedInheritedEventCount ?? 0)
595 if (!this.header.isSeeded && inheritedEventCount !== 0) {
596 throw new Error('unseeded session inherited event count must be 0')
597 }
598 if (inheritedEventCount > this.log.length) {
599 throw new Error('session inherited event count exceeds its event log')
600 }
601 const seedMarker = this.log[inheritedEventCount]
602 const markedSeed = seedMarker?.type === 'session/end-seed' && seedMarker.data.inherited === true
603 if (mode === 'snapshot' && this.header.isSeeded && inheritedEventCount !== this.log.length && !markedSeed) {
604 throw new Error('seeded session constructor seed must equal its inherited prefix or mark its inherited cut')
605 }
606 if (markedSeed && this.log.slice(inheritedEventCount + 1).some(event => event.type === 'session/end-seed' && event.data.inherited === true)) {
607 throw new Error('session inherited event count must identify the final inherited marker')
608 }
609 this.inheritedEventCount = inheritedEventCount
610 this.firstLifecycleSeq = mode === 'snapshot' && this.header.isSeeded ? inheritedEventCount : this.firstLiveSeq
611 // A fresh seeded child always owns one tagged marker at its inherited cut,
612 // even when the copied prefix already ends in an ancestor marker. Restore
613 // retains that durable marker and appends only the ordinary resume marker.
614 if (seed !== undefined && mode === 'snapshot' && this.header.isSeeded && !markedSeed) {
615 this.append('session/end-seed', { inherited: true })
616 } else if (seed !== undefined && !(mode === 'snapshot' && this.header.isSeeded) && this.log.at(-1)?.type !== 'session/end-seed') {
617 this.append('session/end-seed', {})
618 }
619 }
620
621 /** Cached immutable full snapshot of the private append-only log. */
622 private eventsSnapshot: readonly SessionEvent[] | undefined
623
624 /**
625 * Return the immutable event stored at one exact sequence number.
626 * @deprecated Existing logic may remain unmigrated for now, but new calls are prohibited.
627 * See the [Agent Note](../../../../.agents/notes/implemented/architecture/2026-09-09-deprecate-synchronous-session-event-reads.md).
628 * @param seq - event sequence number.
629 * @returns the accepted event, or undefined when the log does not contain it.
630 */
631 eventAt(seq: SessionSeq): SessionEvent | undefined {
632 return this.log[seq]
633 }
634
635 /**
636 * Materialize an immutable snapshot of a half-open event sequence range.
637 * A full current snapshot is reused until the next append; every previously
638 * returned snapshot remains stable after later appends.
639 * @deprecated Existing logic may remain unmigrated for now, but new calls are prohibited.
640 * See the [Agent Note](../../../../.agents/notes/implemented/architecture/2026-09-09-deprecate-synchronous-session-event-reads.md).
641 * @param fromSeq - non-negative inclusive sequence number; defaults to the log start.
642 * @param toSeqExclusive - non-negative exclusive sequence number; defaults to the current end.
643 * @returns a frozen array of the selected deeply frozen events.
644 */
645 snapshotEvents(
646 fromSeq: SessionLogOffset = SessionLogOffset(0),
647 toSeqExclusive: SessionLogOffset = this.seq,
648 ): readonly SessionEvent[] {
649 if (fromSeq === 0 && toSeqExclusive === this.log.length) {
650 this.eventsSnapshot ??= Object.freeze([...this.log])
651 return this.eventsSnapshot
652 }
653 return Object.freeze(this.log.slice(fromSeq, toSeqExclusive))
654 }
655
656 /**
657 * Return this Session's events after its fork-inherited prefix.
658 * @deprecated Existing logic may remain unmigrated for now, but new calls are prohibited.
659 * See the [Agent Note](../../../../.agents/notes/implemented/architecture/2026-09-09-deprecate-synchronous-session-event-reads.md).
660 * @returns a fresh array containing child-owned events in log order.
661 */
662 ownEvents(): readonly SessionEvent[] {
663 // oxlint-disable-next-line typescript/no-deprecated -- Deprecated reader delegates to the deprecated range read.
664 return this.snapshotEvents(this.inheritedEventCount)
665 }
666
667 /**
668 * Whether one existing event position is outside the fork-inherited prefix.
669 * @param seq - event position in this Session.
670 * @returns true when the event belongs to this Session rather than its parent.
671 */
672 isOwnSeq(seq: SessionSeq): boolean {
673 return seq >= this.inheritedEventCount && seq < this.seq
674 }
675
676 /** The next event's sequence number — always the log length (the `seq = log.length` contiguity contract). */
677 get seq(): SessionLogOffset {
678 return SessionLogOffset(this.log.length)
679 }
680
681 /**
682 * Append one typed event to the log and synchronously notify observers via
683 * the store-owned, module-private publication hooks. The hot path never blocks
684 * on I/O — persistence plugins buffer asynchronously. Once the event enters
685 * the log, the append is committed: observer failures are logged and
686 * contained per listener, so they do not change the return value or prevent
687 * later listeners from observing the same accepted event.
688 *
689 * @param type - The event type (key of {@link SessionEventMap}).
690 * @param data - The event payload; must be JSON-serializable.
691 * @param opts - Surface metadata: `surfaceOp` controls how the event enters
692 * the ordered surface; `sourceEventSeqs` lists the seq numbers of earlier
693 * events this one derives from. REQUIRED for
694 * {@link SurfaceEventType} events (every message-producing event must
695 * declare how it joins the surface, the sole source of derived model
696 * history) and
697 * rejected by the compiler for non-surface types like `turn/start` or
698 * `assistant/attempt`. Assistant messages embed their exact provider
699 * stream and cannot cite top-level source events.
700 * @returns the logged event — its assigned `seq`/`time` plus the SNAPSHOT of
701 * `data` that entered the log, so reading `event.data` back sees the logged
702 * value, never the caller's still-mutable input.
703 * @throws if `data` or surface metadata is not losslessly JSON-serializable
704 * (BigInt, function, symbol, undefined, negative zero, non-finite number,
705 * circular reference, sparse array, or an exotic object such as
706 * Map/Set/Date/class instance), or when the candidate violates the
707 * request-header empty-field or tool-error consistency rules, or the
708 * canonical surface contract (marker shape and eligibility, unique
709 * earlier source-event references, positional replacement validity, and complete
710 * shadowed-node coverage). One iterative pass reads, validates, and
711 * copies each nested value once, so a stateful getter cannot supply one value
712 * to validation and another to storage. The event log is the durable source
713 * of truth, so a bad event fails at the append site rather than later during
714 * a backend flush. A synchronous internal dispatch validation failure or an
715 * append reentered while this acceptance/publication boundary is open also
716 * rejects before the log changes.
717 */
718 append<T extends SessionEventType>(
719 type: T,
720 data: SessionEventMap[T],
721 ...opts: T extends SurfaceEventType ? [opts: SurfaceIntent<T>] : []
722 ): SessionEvent<T> {
723 const surfaceOpts: SurfaceIntent | undefined = opts[0]
724 const surfaceMetadata = {
725 ...surfaceOpts?.sourceEventSeqs === undefined ? {} : { sourceEventSeqs: surfaceOpts.sourceEventSeqs },
726 ...surfaceOpts?.surfaceOp === undefined ? {} : { surfaceOp: surfaceOpts.surfaceOp },
727 }
728 const dataSnapshot = snapshotJsonValue(data)
729 if (dataSnapshot === undefined) {
730 throw new Error(`session event "${type}" carries non-JSON-serializable data`)
731 }
732 const surfaceMetadataSnapshot = snapshotJsonValue(surfaceMetadata)
733 if (surfaceMetadataSnapshot === undefined) {
734 throw new Error(`session event "${type}" carries non-JSON-serializable surface metadata`)
735 }
736 const entry = attachments.get(this)
737 if (entry?.appending) {
738 throw new Error('session append cannot reenter while another append is being published')
739 }
740 const event = deepFreeze({
741 type,
742 seq: SessionSeq(this.log.length),
743 time: Date.now(),
744 data: dataSnapshot,
745 ...(surfaceMetadataSnapshot as { surfaceOp?: unknown; sourceEventSeqs?: unknown }),
746 } as unknown as SessionEvent<T>)
747 validateSessionEventData(event, `session event "${type}" at seq ${event.seq}`)
748 this.surfaceManager.validateNext(event as SessionEvent)
749
750 if (entry !== undefined) entry.appending = true
751 try {
752 let callbacks: SessionCallback[] | undefined
753 const callbackArgs: unknown[] = [this, event]
754 if (entry !== undefined) {
755 callbacks = collectSessionCallbacks(entry.emitCtx, [entry.carrier, 'session/event', ...callbackArgs])
756 }
757 this.log.push(event as SessionEvent)
758 this.eventsSnapshot = undefined
759 if (callbacks !== undefined && entry !== undefined) {
760 invokeContainedSessionObservers(entry.emitCtx, 'session/event', entry.id, callbackArgs, callbacks)
761 }
762 return event
763 } finally {
764 if (entry !== undefined) {
765 entry.appending = false
766 if (entry.detachRequested && !entry.announcing) entry.detach()
767 }
768 }
769 }
770
771 /** Cached fold of the request-header events — see {@link requestHeader}. */
772 private headerFold: EpochHeader | undefined
773 /** Log position (events consumed) the header fold has reached. */
774 private headerFoldSeq = 0
775
776 /**
777 * The {@link EpochHeader} in force after the log's last header event — the
778 * header the NEXT request will be compared against — or undefined before
779 * the first `request/header` snapshot. The live, incrementally-maintained
780 * form of `foldRequestHeader(session.snapshotEvents())`: each header event is folded
781 * once, when first seen, so a per-step read costs O(new events).
782 * @returns the folded header, or undefined when no header event exists yet.
783 */
784 requestHeader(): EpochHeader | undefined {
785 if (this.headerFoldSeq < this.log.length) {
786 // Frozen on update: the fold is session state exposed by reference — a
787 // consumer mutating it in place (instead of building a replacement)
788 // would desync every later comparison against the log, so mutation
789 // throws instead.
790 this.headerFold = deepFreeze(foldRequestHeader(this.log.slice(this.headerFoldSeq), this.headerFold))
791 this.headerFoldSeq = this.log.length
792 }
793 return this.headerFold
794 }
795
796 /** Cached fold of `request/context` events. */
797 private contextFold: RequestContext | undefined
798 private contextFoldSeq = 0
799
800 /**
801 * Return the latest resolved route metadata, or `undefined` before the first
802 * `request/context` event. Each event is folded once.
803 * @returns the latest immutable route metadata.
804 */
805 requestContext(): RequestContext | undefined {
806 if (this.contextFoldSeq < this.log.length) {
807 for (const event of this.log.slice(this.contextFoldSeq)) {
808 if (event.type === 'request/context') this.contextFold = deepFreeze({ ...event.data })
809 }
810 this.contextFoldSeq = this.log.length
811 }
812 return this.contextFold
813 }
814
815 /** Cached historical tool definitions and updates for request projection. */
816 private readonly toolHistoryProjection = new ToolHistoryProjection()
817 /** Index of the next committed event not yet consumed by the tool-history fold. */
818 private toolHistorySeq = 0
819
820 /**
821 * Fold unseen committed events into capability-independent tool history.
822 * Initial access reconstructs inherited history; later reads consume only new events.
823 * @returns an immutable snapshot for LLM request projection, including historical addition definitions.
824 */
825 toolHistory(): ToolHistory {
826 for (const event of this.log.slice(this.toolHistorySeq)) this.toolHistoryProjection.apply(event)
827 this.toolHistorySeq = this.log.length
828 return this.toolHistoryProjection.snapshot()
829 }
830
831 /** The derived-message cache: frozen projections, extended per unseen node. */
832 private derived: Message[] = []
833 /** Surface position (nodes projected) the cache has reached. */
834 private derivedNodes = 0
835 /** {@link SurfaceManager.contentGeneration} the cache was built under. */
836 private derivedGeneration = 0
837
838 /**
839 * Derive the LLM message history by walking the ordered sequences of
840 * message-producing events maintained by `surfaceOp` markers. The
841 * surface is the single source of derived history: every message-producing
842 * append records its `surfaceOp`, so a raw event with no marker (a chunk, a
843 * turn boundary) is correctly absent, and a compaction `replace` deletes the
844 * shadowed nodes from the derivation. The projection rules are
845 * {@link deriveEventMessage}, with logged message projections applied
846 * without changing node membership or message identity.
847 *
848 * CACHED: pure tail growth costs O(new nodes); a replacement or message projection
849 * ({@link SessionSurface.contentGeneration}) rebuilds. The returned array is
850 * a fresh snapshot per call (later appends never grow an array a caller
851 * already holds); the `Message` objects in it are SHARED and **deep-frozen**.
852 * Unchanged content reuses frozen event data; projected blocks are frozen
853 * derived copies. Consumers cannot mutate the log through either form.
854 * @returns a fresh array of the shared, frozen derived history.
855 */
856 deriveMessages(): Message[] {
857 const surface = this.surface
858 const nodes = surface.nodes
859 const generation = surface.contentGeneration
860 if (generation !== this.derivedGeneration) {
861 this.derived = []
862 this.derivedNodes = 0
863 this.derivedGeneration = generation
864 }
865 for (const seq of nodes.slice(this.derivedNodes)) {
866 // Surface sequences are built from this.log — seq is always a valid
867 // index by construction. The non-null assertion expresses that invariant.
868 // oxlint-disable-next-line typescript/no-non-null-assertion
869 const msg = this.deriveEventMessage(this.log[seq]!)
870 // A surface node is one of the five message-producing types, but an
871 // empty-content assistant/message (a max-tokens step that hosts only
872 // usage) derives to null and must not enter the transcript.
873 if (msg) this.derived.push(msg)
874 }
875 this.derivedNodes = nodes.length
876 return [...this.derived]
877 }
878
879 /**
880 * Project one event with all committed message projections applied.
881 * The original durable event remains unchanged.
882 * @param event - the event to project.
883 * @returns the derived message, or null when the event produces none.
884 */
885 deriveEventMessage(event: SessionEvent): Message | null {
886 return this.surfaceManager.deriveEventMessage(event)
887 }
888}
889
890/** A fork source: either the live session object or its live store id. */
891export type SessionForkSource = Session | SessionId
892
893/**
894 * Rejection codes for session forking: the fork source id is unknown to the
895 * live store (`SESSION_NOT_FOUND`) or names a session object that is not the
896 * store's live instance (`SESSION_NOT_LIVE`); the requested child id is
897 * already taken (`SESSION_ALREADY_EXISTS`); or the boundary is not a contiguous
898 * existing seq (`INVALID_BOUNDARY`).
899 */
900export type SessionForkErrorCode =
901 | 'SESSION_NOT_FOUND'
902 | 'SESSION_NOT_LIVE'
903 | 'SESSION_ALREADY_EXISTS'
904 | 'INVALID_BOUNDARY'
905
906/** Typed error for session fork rejections. */
907export class SessionForkError extends Error {
908 constructor(message: string, public readonly code: SessionForkErrorCode) {
909 super(message)
910 this.name = 'SessionForkError'
911 }
912}
913
914/**
915 * In-memory session store (`ctx.sessions`).
916 *
917 * Persistence is intentionally not implemented here — the agent lifecycle
918 * attaches a session-log writer to each published session's write handle;
919 * a session published outside that lifecycle persists nothing.
920 */
921export class SessionStore extends Service {
922 private store = new Map<SessionId, SessionEntry>()
923 private counter = 0
924 private readonly projections: SessionMessageProjection[] = []
925
926 /** Borrowed definitions for detached replay; contributions live until their registering fibers unload. */
927 get messageProjections(): readonly SessionMessageProjection[] {
928 return this.projections
929 }
930
931 /**
932 * Register one event interpreter for live creation, restore, and fork.
933 * Disposing the contribution makes sessions that used it refuse further derivation.
934 * @param projection - pure definition owned by the event's plugin.
935 * @returns the fiber-owned disposer.
936 * @throws when another definition already owns this event type.
937 */
938 registerMessageProjection(projection: SessionMessageProjection): () => Promise<void> {
939 if (this.projections.some(item => item.type === projection.type)) {
940 throw new Error(`session message projection "${projection.type}" is already registered`)
941 }
942 return this.ctx.effect(() => {
943 this.projections.push(projection)
944 return () => { this.projections.splice(this.projections.indexOf(projection), 1) }
945 }, 'sessions.registerMessageProjection()')
946 }
947
948 constructor(ctx: Context) {
949 super(ctx, 'sessions')
950 ctx.inject(['typert'], (typeCtx) => {
951 typeCtx.typert.lookups.register('session', {
952 parameter: 'session',
953 wire: 'sessionId',
954 hostTypeSymbol: '@deepseek-ai/dsh-session#Session',
955 wireTypeSymbol: '@deepseek-ai/dsh-session/types#SessionId',
956 resolve: sessionId => this.get(sessionId),
957 })
958 })
959 }
960
961 /**
962 * Create a session owned by the calling fiber: disposing that fiber stops
963 * event notification and removes the session from the store. `options.seed`
964 * populates the session with a copy of those events (replay/fork);
965 * `options.meta` attaches creation metadata (validated absolute `cwd`, seed
966 * and parent lineage, and delegation depth) as the immutable
967 * {@link SessionHeader} (the store fills `version`/`id`/`createdAt`).
968 *
969 * For an agent whose session must be torn down IN ORDER with its loop (so the
970 * loop's final events are published before the store attachment ends), do NOT use this
971 * — fold the session lifecycle into the agent's own effect via
972 * {@link prepare} + {@link enter} + {@link announce} (see
973 * `dsh-agent-loop`'s creation transaction).
974 *
975 * @param id - the session id; omitted, the store mints `session-<n>`.
976 * @param options - seed events and/or creation metadata for the header.
977 * @returns the live session, already entered and announced.
978 * @throws if a session with `id` already exists, metadata is not a plain
979 * lossless-JSON record with valid scalar fields, or `meta.cwd` is a
980 * non-absolute path (storage backends key directories off it).
981 */
982 create(id?: SessionId, options?: CreateSessionOptions): Session {
983 const session = this.prepare(id, options)
984 // Single effect owned by the calling fiber. Yield the detach BEFORE
985 // announcing so a throwing `session/created` listener rolls the attach back
986 // (the generator effect disposes already-yielded disposers on a throw)
987 // instead of leaking the store entry and its publication hooks.
988 this.ctx.effect(function* (this: SessionStore) {
989 yield this.enter(session)
990 this.announce(session)
991 }.bind(this), 'sessions.create()')
992 return session
993 }
994
995 /**
996 * Build a session WITHOUT entering it into the store — validate the id/cwd and
997 * construct the {@link Session} (with its immutable {@link SessionHeader}).
998 * Pairs with {@link enter} + {@link announce}: a caller that owns a composite
999 * `ctx.effect` (the agent factory) folds the session lifecycle into that ONE
1000 * effect so a fiber unload tears the session + agent down as a single ORDERED
1001 * chain rather than as racing sibling effects — which would remove the publication hooks
1002 * before the driver's closing events commit, dropping them.
1003 *
1004 * @param id - the session id; omitted, the store mints `session-<n>`.
1005 * @param options - seed events and/or creation metadata for the header. With
1006 * `eventState`, every seed event is either independently owned or any
1007 * shared value is deeply frozen; {@link Session.fromRestore} validates and
1008 * adopts those values without copying or freezing them.
1009 * @returns the constructed session, NOT yet in the store.
1010 * @throws if a session with `id` already exists, metadata is not a plain
1011 * lossless-JSON record with valid scalar fields, or `meta.cwd` is a
1012 * non-absolute path.
1013 */
1014 prepare(id?: SessionId, options?: PrepareSessionOptions): Session {
1015 let sessionId: SessionId
1016 if (id === undefined) {
1017 do sessionId = brandString<SessionId>(`session-${++this.counter}`)
1018 while (this.store.has(sessionId))
1019 } else {
1020 sessionId = brandString<SessionId>(id)
1021 }
1022 if (this.store.has(sessionId)) throw new Error(`session "${sessionId}" already exists`)
1023 if (options !== undefined) {
1024 const { eventState } = options
1025 switch (eventState) {
1026 case 'detached':
1027 case 'shared-frozen':
1028 return Session.fromRestore(
1029 sessionId,
1030 options.seed,
1031 options.meta,
1032 options.inheritedEventCount,
1033 eventState,
1034 this.projections,
1035 )
1036 case undefined:
1037 break
1038 /* v8 ignore next -- closed-union exhaustiveness guard */
1039 default:
1040 assertNever(eventState, 'SessionStore.prepare event state')
1041 }
1042 }
1043 const seed = options?.seed
1044 const meta = options?.meta
1045 const header: SessionHeader = {
1046 version: SESSION_FORMAT_VERSION,
1047 id: sessionId,
1048 createdAt: meta?.createdAt ?? Date.now(),
1049 ...meta?.cwd === undefined ? {} : { cwd: meta.cwd },
1050 ...meta?.parentSession === undefined ? {} : { parentSession: meta.parentSession },
1051 isSeeded: meta?.isSeeded ?? false,
1052 ...meta?.origin === undefined ? {} : { origin: meta.origin },
1053 ...meta?.delegationDepth === undefined ? {} : { delegationDepth: meta.delegationDepth },
1054 ...meta?.agentPreset === undefined ? {} : { agentPreset: meta.agentPreset },
1055 }
1056 return Session.create(sessionId, seed, header, options?.inheritedEventCount, this.projections)
1057 }
1058
1059 /**
1060 * Enter a {@link prepare}d session into the store: install the module-private
1061 * append publication hooks and add it to the store. Returns the DETACH
1062 * disposer (hooks + store removal). Does NOT emit `session/created` —
1063 * the caller yields this disposer inside its effect and THEN calls
1064 * {@link announce}, so a throwing `session/created` listener rolls the attach
1065 * back instead of leaking it.
1066 *
1067 * Re-checks the id for a duplicate: `prepare` and `enter` are public
1068 * cross-package primitives and a caller may interleave arbitrary work (or
1069 * another create) between them, so a stale prepared session must NOT overwrite
1070 * a live store entry of the same id — its detach disposer would later delete
1071 * the REAL session. The {@link create} convenience and the agent factory call
1072 * the two back-to-back so they never trip this, but the public API cannot
1073 * assume that.
1074 *
1075 * @param session - a {@link prepare}d session not yet in the store.
1076 * @returns the detach disposer (publication hooks + store removal). When called from
1077 * a synchronous `session/created` listener, removal and disposal wait until
1078 * that creation dispatch unwinds.
1079 * @throws if a session with this id is already in the store.
1080 */
1081 enter(session: Session): () => void {
1082 const id = session.id
1083 const carrier = scopeTarget(session, scopeOf(this.ctx))
1084 // This is the authoritative collision boundary after arbitrary unpublished
1085 // preparation. Only one exact same-id transaction can publish.
1086 if (this.store.has(id)) throw new Error(`session "${id}" already exists`)
1087 if (attachments.has(session)) throw new Error(`session "${id}" is already attached to a store`)
1088 const entry: SessionEntry = {
1089 id,
1090 session,
1091 carrier,
1092 emitCtx: this.ctx,
1093 announced: false,
1094 announcing: false,
1095 appending: false,
1096 detachRequested: false,
1097 detach: () => { this.detachEntered(entry) },
1098 }
1099 this.store.set(id, entry)
1100 attachments.set(session, entry)
1101 let entered = true
1102 const detach = (): void => {
1103 if (!entered) return
1104 entered = false
1105 // A lifecycle listener may own the advanced detach capability. Keep the
1106 // entry and its publication hooks live until synchronous creation or append
1107 // publication unwinds, then publish the paired disposal edge.
1108 if (entry.announcing || entry.appending) {
1109 entry.detachRequested = true
1110 return
1111 }
1112 entry.detach()
1113 }
1114 return detach
1115 }
1116
1117 /** Remove one exact entered session and emit its paired disposal when announced. */
1118 private detachEntered(entry: SessionEntry): void {
1119 entry.detachRequested = false
1120 // A stale capability cannot remove observers or storage belonging to a
1121 // later same-id lifecycle.
1122 /* v8 ignore next -- enter() rejects replacement while this single-shot detach capability is live. */
1123 if (this.store.get(entry.id) !== entry) return
1124 this.store.delete(entry.id)
1125 attachments.delete(entry.session)
1126 if (entry.announced) this.emitDisposed(entry)
1127 }
1128
1129 /** Emit `session/created` exactly once for an {@link enter}ed session (with
1130 * the carrier {@link enter} captured). Separate from {@link enter} so the
1131 * caller can yield the detach disposer first (rollback safety — see
1132 * {@link enter}).
1133 * @param session - the entered session to announce to listeners.
1134 * @throws if the session is not live or its announcement already began,
1135 * including a reentrant call from a creation listener. */
1136 announce(session: Session): void {
1137 const entry = this.liveEntryFor(session)
1138 if (entry.announced || entry.announcing) {
1139 throw new Error(`session "${entry.id}" was already announced`)
1140 }
1141 // Mark before emit: Cordis emit may deliver to earlier listeners and then
1142 // throw. Rollback must still pair that partial creation with disposal, and
1143 // a listener cannot recursively create a second lifecycle edge.
1144 entry.announced = true
1145 const callbackArgs: unknown[] = [session]
1146 entry.announcing = true
1147 try {
1148 const callbacks = collectSessionCallbacks(this.ctx, [entry.carrier, 'session/created', session])
1149 for (const callback of callbacks) {
1150 // Synchronous throws intentionally propagate and veto publication; the
1151 // yielded detach then emits the paired disposal edge. An async function
1152 // is nevertheless assignable to a void listener, so observe its returned
1153 // promise: rejection is too late to roll back and must be logged instead
1154 // of becoming unhandled.
1155 const returned: unknown = callback(...callbackArgs)
1156 void Promise.resolve(returned).catch((error: unknown) => {
1157 this.ctx.logger.warn(`session "${entry.id}": session/created listener rejected: ${String(error)}`)
1158 })
1159 }
1160 } finally {
1161 entry.announcing = false
1162 if (entry.detachRequested && !entry.appending) entry.detach()
1163 }
1164 }
1165
1166 /** Emit the paired teardown notification with per-listener containment. */
1167 private emitDisposed(entry: SessionEntry): void {
1168 const callbackArgs: unknown[] = [entry.session]
1169 try {
1170 const callbacks = collectSessionCallbacks(this.ctx, [entry.carrier, 'session/disposed', entry.session])
1171 invokeContainedSessionObservers(this.ctx, 'session/disposed', entry.id, callbackArgs, callbacks)
1172 } catch (error: unknown) {
1173 this.ctx.logger.warn(`session "${entry.id}": session/disposed dispatch threw: ${String(error)}`)
1174 }
1175 }
1176
1177 /**
1178 * Dispatch the awaited `session/flush` durability checkpoint for `session`,
1179 * with the carrier captured at {@link enter}. THE flush entry point: the
1180 * store owns the carrier, so callers (the checkpoint policy's per-request
1181 * barrier, goal-round-driver's idle checkpoint, teardown drains, and consumers
1182 * that flush themselves before reading storage) must come through here
1183 * rather than dispatch a raw `ctx.parallel('session/flush', …)` — one owner
1184 * and one spelling.
1185 * @param session - the session whose buffered events must reach durable storage.
1186 * @returns whether at least one durability listener participated, after every
1187 * listener has settled successfully.
1188 * @throws the first registered listener failure after every listener settles.
1189 */
1190 async flush(session: Session): Promise<boolean> {
1191 const { carrier } = this.liveEntryFor(session)
1192 const callbackArgs: unknown[] = [session]
1193 const callbacks = collectSessionCallbacks(this.ctx, [carrier, 'session/flush', session])
1194 const results = await Promise.allSettled(callbacks.map((callback) => {
1195 try {
1196 return callback(...callbackArgs)
1197 } catch (error: unknown) {
1198 // Preserve the listener's exact rejection value; flush is a caller-owned
1199 // failure boundary, and Cordis listeners may throw arbitrary values.
1200 // oxlint-disable-next-line typescript/prefer-promise-reject-errors
1201 return Promise.reject(error)
1202 }
1203 }))
1204 const failure = results.find((result): result is PromiseRejectedResult => result.status === 'rejected')
1205 if (failure !== undefined) throw failure.reason
1206 return callbacks.length > 0
1207 }
1208
1209 /** Return the exact live entry; detached/prepared objects reject. */
1210 private liveEntryFor(session: Session): SessionEntry {
1211 const entry = attachments.get(session)
1212 if (entry === undefined || this.store.get(entry.id) !== entry) {
1213 throw new Error(`session "${session.id}" is not live in this store`)
1214 }
1215 return entry
1216 }
1217
1218 /**
1219 * Look up a live session.
1220 * @param id - the session id to look up.
1221 * @returns the session, or undefined when no live session has that id.
1222 */
1223 get(id: SessionId): Session | undefined {
1224 return this.store.get(id)?.session
1225 }
1226
1227 /**
1228 * All live sessions, in creation order.
1229 * @returns a fresh array; mutating it does not affect the store.
1230 */
1231 list(): Session[] {
1232 return [...this.store.values()].map(entry => entry.session)
1233 }
1234
1235 /**
1236 * Create a live child session from an exact prefix of a live source.
1237 * `boundary` is an inclusive source event seq; omitted means the source's
1238 * current last event. An open tail receives synthetic tool results and
1239 * step/turn closers with the forked cause. Closed steps and turns remain
1240 * unchanged, including any failed tool calls already missing results.
1241 * `inheritedEventCount` counts only copied source events, excluding these closers.
1242 *
1243 * @param source - Live source session object or id.
1244 * @param boundary - Inclusive source event seq to fork through; omitted means
1245 * the source's current last event, and omitted on an empty source forks an
1246 * empty child.
1247 * @param childSessionId - Optional child session id; omitted delegates to
1248 * `SessionStore`'s id policy.
1249 * @returns The created live child session.
1250 */
1251 fork(source: SessionForkSource, boundary?: SessionSeq, childSessionId?: SessionId): Session {
1252 if (childSessionId !== undefined && this.get(childSessionId) !== undefined) {
1253 throw new SessionForkError(`session "${childSessionId}" already exists`, 'SESSION_ALREADY_EXISTS')
1254 }
1255 const liveSource = this._resolveForkSource(source)
1256 // oxlint-disable-next-line typescript/no-deprecated -- Existing fork snapshot read; migration deferred.
1257 const events = liveSource.snapshotEvents()
1258 const resolved = this._forkBoundary(liveSource.id, events, boundary)
1259 const seed = resolved === undefined ? [] : buildForkSeed(events, resolved)
1260 return this.create(childSessionId, {
1261 seed,
1262 inheritedEventCount: SessionLogOffset(resolved === undefined ? 0 : resolved + 1),
1263 meta: {
1264 ...liveSource.header.cwd !== undefined ? { cwd: liveSource.header.cwd } : {},
1265 parentSession: liveSource.id,
1266 isSeeded: true,
1267 },
1268 })
1269 }
1270
1271 private _forkBoundary(
1272 sessionId: SessionId, events: readonly SessionEvent[], requestedBoundary: SessionSeq | undefined,
1273 ): SessionSeq | undefined {
1274 const lastEvent = events.at(-1)
1275 let boundary: SessionSeq
1276 if (requestedBoundary !== undefined) {
1277 boundary = requestedBoundary
1278 } else {
1279 if (lastEvent === undefined) return undefined
1280 boundary = lastEvent.seq
1281 }
1282 if (!Number.isSafeInteger(boundary) || boundary < 0) {
1283 throw new SessionForkError(
1284 `fork boundary for session "${sessionId}" must be a non-negative safe integer, got ${String(boundary)}`,
1285 'INVALID_BOUNDARY',
1286 )
1287 }
1288 if (boundary >= events.length) {
1289 const lastSeq = lastEvent?.seq
1290 throw new SessionForkError(
1291 `fork boundary ${boundary} does not exist in session "${sessionId}" (last seq: ${lastSeq ?? 'none'})`,
1292 'INVALID_BOUNDARY',
1293 )
1294 }
1295
1296 const boundaryEvent = events[boundary]
1297 if (boundaryEvent === undefined || boundaryEvent.seq !== boundary) {
1298 throw new SessionForkError(
1299 `fork boundary ${boundary} does not match a contiguous event seq in session "${sessionId}"`,
1300 'INVALID_BOUNDARY',
1301 )
1302 }
1303 return boundary
1304 }
1305
1306 private _resolveForkSource(source: SessionForkSource): Session {
1307 if (typeof source === 'string') {
1308 const session = this.get(source)
1309 if (session === undefined) throw new SessionForkError(`session "${source}" not found`, 'SESSION_NOT_FOUND')
1310 return session
1311 }
1312
1313 const live = this.get(source.id)
1314 if (live === undefined) {
1315 throw new SessionForkError(`session "${source.id}" not found`, 'SESSION_NOT_FOUND')
1316 }
1317 if (live !== source) throw new SessionForkError(`session "${source.id}" is not the live store instance`, 'SESSION_NOT_LIVE')
1318 return source
1319 }
1320
1321}
1322
1323export { decodeSeqRanges, encodeSeqRanges } from './seq-ranges.ts'
1324export default SessionStore