index.ts 47 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038
  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. import { Context, Service } from 'cordis'
  9. import { isAbsolute } from 'node:path'
  10. import { deepFreeze } from '@deepseek-ai/dsh-llm'
  11. import { scopeOf, scopeTarget } from '@deepseek-ai/dsh-scope'
  12. import type { Scoped } from '@deepseek-ai/dsh-scope'
  13. import type { Message } from '@deepseek-ai/dsh-llm'
  14. import { SESSION_FORMAT_VERSION, SessionId } from './types.ts'
  15. import type { CreateSessionOptions, EpochHeader, SessionEvent, SessionEventMap, SessionEventType, SessionHeader, SurfaceIntent, SurfaceEventType } from './types.ts'
  16. import { snapshotJsonValue } from './json.ts'
  17. import { SurfaceManager } from './surface.ts'
  18. import type { SessionSurface } from './surface.ts'
  19. import { foldRequestHeader } from './request-header.ts'
  20. export * from './types.ts'
  21. export type { AssistantMessage, ToolResultMessage, UserMessage } from '@deepseek-ai/dsh-llm'
  22. export { isJsonValue, snapshotJsonValue } from './json.ts'
  23. export type { JsonValue } from './json.ts'
  24. export { interruptedTurnClosers, TOOL_NOT_STARTED, TOOL_OUTCOME_UNKNOWN } from './repair.ts'
  25. export { decodeStorageRecord, packChunkRuns } from './chunk-rows.ts'
  26. export type { ChunkRow, StorageRecord } from './chunk-rows.ts'
  27. export type { SessionSurface, SurfaceFoldReplacement, SurfaceFoldResult } from './surface.ts'
  28. export { foldSurface, isSurfaceEvent, isSurfaceEligibleType } from './surface.ts'
  29. export { canonicalHeader, foldRequestHeader, headerEquals } from './request-header.ts'
  30. /**
  31. * Find the latest closed message-triggered turn, ignoring other triggers and
  32. * between-turn events.
  33. * @param events - session events, or an owned suffix, to inspect.
  34. * @returns the latest matching turn end, or `undefined`.
  35. */
  36. export function findLastMessageTurnEnd(
  37. events: readonly SessionEvent[],
  38. ): SessionEvent<'turn/end'> | undefined {
  39. const messageTurns = new Set<number>()
  40. let latest: SessionEvent<'turn/end'> | undefined
  41. for (const event of events) {
  42. if (event.type === 'turn/start') {
  43. if (event.data.trigger.kind === 'message') messageTurns.add(event.data.turn)
  44. continue
  45. }
  46. if (event.type === 'turn/end' && messageTurns.delete(event.data.turn)) latest = event
  47. }
  48. return latest
  49. }
  50. declare module 'cordis' {
  51. interface Context {
  52. sessions: SessionStore
  53. }
  54. interface Events {
  55. /**
  56. * Creation announcement during session publication. A synchronous throw vetoes and rolls
  57. * back with a paired disposal; detach requested during dispatch is deferred.
  58. * A returned-promise rejection is logged but cannot retroactively veto this
  59. * synchronous boundary.
  60. * Scope-filtered dispatch (`@deepseek-ai/dsh-scope`): agent-scoped listeners
  61. * receive only sessions entered through that agent's context.
  62. * @param session - the session just entered and announced.
  63. * @dshScopeScan unsupported
  64. * @mode emit
  65. */
  66. 'session/created'(this: Scoped<Session>, session: Session): void
  67. /**
  68. * Emitted once when an announced session leaves the store, including
  69. * publication rollback, but never for an entry whose creation announcement
  70. * did not begin. Listener failures are logged and contained.
  71. * Scope-filtered dispatch (`@deepseek-ai/dsh-scope`) reuses the owner scope.
  72. * @param session - the session that is no longer live in the store.
  73. * @dshScopeScan unsupported
  74. * @mode emit
  75. */
  76. 'session/disposed'(this: Scoped<Session>, session: Session): void
  77. /**
  78. * Post-commit, fire-and-forget append feed. The listener snapshot resolves
  79. * before the log push, but callbacks run after it; observer failures are
  80. * logged and contained without making the committed append fail.
  81. * Scope-filtered dispatch (`@deepseek-ai/dsh-scope`): agent-scoped listeners
  82. * receive only events from sessions entered through that agent's context.
  83. * @param session - the session whose log grew.
  84. * @param event - the appended event, exactly as recorded.
  85. * @dshScopeScan unsupported
  86. * @mode emit
  87. */
  88. 'session/event'(this: Scoped<Session>, session: Session, event: SessionEvent): void
  89. /**
  90. * Awaited parallel durability checkpoint: every listener runs and the
  91. * caller awaits all of them, with no waterfall veto. Dispatch through
  92. * {@link SessionStore.flush}. Scope-filtered dispatch
  93. * (`@deepseek-ai/dsh-scope`) reuses the session's owner scope.
  94. * @param session - the session whose buffered events must reach durable storage.
  95. * @dshScopeScan unsupported
  96. * @mode parallel
  97. */
  98. 'session/flush'(this: Scoped<Session>, session: Session): Promise<void> | void
  99. }
  100. }
  101. /** Detach, validate, and freeze the creation metadata published by a session. */
  102. function snapshotSessionHeader(id: SessionId, source?: SessionHeader): SessionHeader {
  103. const input: unknown = source === undefined
  104. ? { version: SESSION_FORMAT_VERSION, id, createdAt: Date.now() }
  105. : source
  106. const snapshot = snapshotJsonValue(input)
  107. if (snapshot === undefined) throw new Error('session header is not losslessly JSON-serializable')
  108. if (snapshot === null || typeof snapshot !== 'object' || Array.isArray(snapshot)) {
  109. throw new Error('session header is not a plain JSON record')
  110. }
  111. const record = snapshot as Record<string, unknown>
  112. if (record.version !== SESSION_FORMAT_VERSION) {
  113. throw new Error(`session header version must be ${SESSION_FORMAT_VERSION}, got ${String(record.version)}`)
  114. }
  115. if (record.id !== id) {
  116. throw new Error(`session header id "${String(record.id)}" does not match session id "${id}"`)
  117. }
  118. if (typeof record.createdAt !== 'number'
  119. || !Number.isSafeInteger(record.createdAt)
  120. || record.createdAt < 0) {
  121. throw new Error('session header createdAt must be a non-negative safe integer')
  122. }
  123. if (record.cwd !== undefined) {
  124. if (typeof record.cwd !== 'string') throw new Error('session header cwd must be a string')
  125. if (!isAbsolute(record.cwd)) {
  126. throw new Error(`session header cwd must be an absolute path, got "${record.cwd}"`)
  127. }
  128. }
  129. if (record.parentSession !== undefined && typeof record.parentSession !== 'string') {
  130. throw new Error('session header parentSession must be a string')
  131. }
  132. if (record.seedLength !== undefined
  133. && (typeof record.seedLength !== 'number' || !Number.isSafeInteger(record.seedLength) || record.seedLength < 0)) {
  134. throw new Error('session header seedLength must be a non-negative safe integer')
  135. }
  136. if (record.delegationDepth !== undefined
  137. && (typeof record.delegationDepth !== 'number' || !Number.isSafeInteger(record.delegationDepth) || record.delegationDepth < 0)) {
  138. throw new Error('session header delegationDepth must be a non-negative safe integer')
  139. }
  140. return deepFreeze(record as unknown as SessionHeader)
  141. }
  142. /**
  143. * Detach one event while preserving deep immutability for its identified message.
  144. * @param event - event imported across a query or persistence boundary.
  145. * @returns a detached event snapshot with a validated, deeply frozen message.
  146. */
  147. export function snapshotSessionEvent<T extends SessionEvent>(event: T): T {
  148. const snapshot = structuredClone(event)
  149. assertMessageEventShape(
  150. snapshot,
  151. `session event at seq ${snapshot.seq}`,
  152. )
  153. switch (snapshot.type) {
  154. case 'user/message':
  155. deepFreeze(snapshot.data)
  156. break
  157. case 'assistant/message':
  158. case 'tool/result':
  159. case 'steering/message':
  160. deepFreeze(snapshot.data.message)
  161. break
  162. default:
  163. // SessionEventMap is merge-extensible; plugin-owned events carry no core message.
  164. break
  165. }
  166. return snapshot
  167. }
  168. /** Validate the fixed event envelope after one-pass JSON materialization. */
  169. function assertSessionEventEnvelope(value: Record<string, unknown>, index: number): asserts value is SessionEvent {
  170. const event = value
  171. if (event['type'] === 'request/header-delta') {
  172. throw new Error(`seed event at index ${index} uses unsupported legacy request/header-delta format`)
  173. }
  174. const allowed = new Set(['type', 'seq', 'time', 'data', 'surfaceOp', 'sourceEventSeqs'])
  175. if (Object.keys(event).some(key => !allowed.has(key))
  176. || !Object.hasOwn(event, 'type') || typeof event['type'] !== 'string'
  177. || !Object.hasOwn(event, 'seq') || typeof event['seq'] !== 'number'
  178. || !Number.isSafeInteger(event['seq']) || event['seq'] < 0
  179. || !Object.hasOwn(event, 'time') || typeof event['time'] !== 'number'
  180. || !Number.isSafeInteger(event['time']) || event['time'] < 0
  181. || !Object.hasOwn(event, 'data')) {
  182. throw new Error(`seed event at index ${index} has an invalid event envelope`)
  183. }
  184. assertCurrentLlmShape(event, index)
  185. assertCurrentTurnEndShape(event, index)
  186. }
  187. /** Reject obsolete request headers and malformed messages at the seed/load boundary. */
  188. function assertCurrentLlmShape(event: Record<string, unknown>, index: number): void {
  189. const data = event['data']
  190. const record = typeof data === 'object' && data !== null
  191. ? data as Record<string, unknown>
  192. : undefined
  193. if (event['type'] === 'request/header') {
  194. const header = record?.['header']
  195. const config = typeof header === 'object' && header !== null ? (header as Record<string, unknown>)['config'] : undefined
  196. if (!hasProviderModel(config)) throw new Error(`seed request/header at index ${index} lacks provider/model`)
  197. const reasoningEffort = (config as Record<string, unknown>)['reasoningEffort']
  198. if (reasoningEffort !== undefined
  199. && (typeof reasoningEffort !== 'string' || reasoningEffort.length === 0)) {
  200. throw new Error(`seed request/header at index ${index} has an invalid reasoningEffort`)
  201. }
  202. }
  203. const type = event['type']
  204. if (type !== 'user/message' && type !== 'assistant/message'
  205. && type !== 'tool/result' && type !== 'steering/message') return
  206. assertMessageEventShape(event, `seed ${type} at index ${index}`)
  207. }
  208. /** Validate only the event-specific invariants needed to safely replay a message. */
  209. function assertMessageEventShape(event: Record<string, unknown>, subject: string): void {
  210. const type = event['type']
  211. if (type !== 'user/message' && type !== 'assistant/message'
  212. && type !== 'tool/result' && type !== 'steering/message') return
  213. const data = event['data']
  214. const record = typeof data === 'object' && data !== null
  215. ? data as Record<string, unknown>
  216. : undefined
  217. const message = type === 'user/message' ? record : record?.['message']
  218. if (typeof message !== 'object' || message === null
  219. || typeof (message as Record<string, unknown>)['id'] !== 'string'
  220. || (message as Record<string, unknown>)['id'] === '') {
  221. throw new Error(`${subject} lacks an identified message`)
  222. }
  223. const messageRecord = message as Record<string, unknown>
  224. const expectedRole = type === 'assistant/message' ? 'assistant' : 'user'
  225. if (messageRecord['role'] !== expectedRole) {
  226. throw new Error(`${subject} message must have role "${expectedRole}"`)
  227. }
  228. const source = messageRecord['source']
  229. if (typeof source !== 'object' || source === null
  230. || typeof (source as Record<string, unknown>)['kind'] !== 'string'
  231. || (source as Record<string, unknown>)['kind'] === '') {
  232. throw new Error(`${subject} message has invalid source`)
  233. }
  234. if (!Array.isArray(messageRecord['content'])) {
  235. throw new Error(`${subject} message has invalid content`)
  236. }
  237. const sourceRecord = source as Record<string, unknown>
  238. if (type === 'assistant/message') {
  239. if (sourceRecord['kind'] !== 'model' || !hasProviderModel(sourceRecord)) {
  240. throw new Error(`${subject} message must have model source`)
  241. }
  242. return
  243. }
  244. if (type !== 'tool/result') return
  245. if (sourceRecord['kind'] !== 'tool'
  246. || typeof sourceRecord['callId'] !== 'string'
  247. || sourceRecord['callId'] === '') {
  248. throw new Error(`${subject} message must have tool source`)
  249. }
  250. const content = messageRecord['content'] as unknown[]
  251. const block = content[0]
  252. if (content.length !== 1 || typeof block !== 'object' || block === null
  253. || (block as Record<string, unknown>)['type'] !== 'tool-result'
  254. || !Array.isArray((block as Record<string, unknown>)['content'])) {
  255. throw new Error(`${subject} message must contain one tool-result block`)
  256. }
  257. if ((block as Record<string, unknown>)['toolCallId'] !== sourceRecord['callId']) {
  258. throw new Error(`${subject} message has mismatched tool call ids`)
  259. }
  260. }
  261. /** Reject legacy aborted outcomes that persisted caller-owned reason detail. */
  262. function assertCurrentTurnEndShape(event: Record<string, unknown>, index: number): void {
  263. if (event['type'] !== 'turn/end') return
  264. const data = event['data']
  265. /* v8 ignore next -- this migration recognizes only the legacy object shape; format-wide payload validation is separate. */
  266. if (typeof data !== 'object' || data === null) return
  267. const reason = (data as Record<string, unknown>)['reason']
  268. /* v8 ignore next -- non-object reasons cannot carry the legacy aborted detail this migration removes. */
  269. if (typeof reason !== 'object' || reason === null || Array.isArray(reason)) return
  270. const record = reason as Record<string, unknown>
  271. if (record['kind'] === 'aborted'
  272. && (Object.keys(record).length !== 1 || !Object.hasOwn(record, 'kind'))) {
  273. throw new Error(`seed turn/end at index ${index} uses unsupported reason-bearing aborted format`)
  274. }
  275. }
  276. /** Whether an unknown value carries the current provider/model pair. */
  277. function hasProviderModel(value: unknown): boolean {
  278. if (typeof value !== 'object' || value === null) return false
  279. const pair = value as Record<string, unknown>
  280. return typeof pair['provider'] === 'string' && pair['provider'].length > 0
  281. && typeof pair['model'] === 'string' && pair['model'].length > 0
  282. }
  283. /** Reject request-header vocabulary removed with the legacy delta codec. */
  284. function assertSupportedRequestHeader(type: string, data: unknown, location: string): void {
  285. if (type === 'request/header-delta') {
  286. throw new Error(`${location} uses unsupported legacy request/header-delta format`)
  287. }
  288. if (type === 'request/header'
  289. && data !== null && typeof data === 'object' && !Array.isArray(data)
  290. && (data as Record<string, unknown>)['reason'] === 'fallback') {
  291. throw new Error(`${location} uses unsupported legacy request/header reason "fallback"`)
  292. }
  293. }
  294. type SessionCallback = (...args: unknown[]) => unknown
  295. /** Resolve one listener snapshot, including Cordis's internal dispatch checks. */
  296. function collectSessionCallbacks(ctx: Context, args: unknown[]): SessionCallback[] {
  297. return [...ctx.events.dispatch('emit', args)] as SessionCallback[]
  298. }
  299. /** Invoke one resolved observe-only listener snapshot with per-listener containment. */
  300. function invokeContainedSessionObservers(
  301. ctx: Context,
  302. name: 'session/event' | 'session/disposed',
  303. id: SessionId,
  304. args: unknown[],
  305. callbacks: SessionCallback[],
  306. ): void {
  307. for (const callback of callbacks) {
  308. try {
  309. const returned: unknown = callback(...args)
  310. void Promise.resolve(returned).catch((error: unknown) => {
  311. ctx.logger.warn(`session "${id}": ${name} listener rejected: ${String(error)}`)
  312. })
  313. } catch (error: unknown) {
  314. ctx.logger.warn(`session "${id}": ${name} listener threw: ${String(error)}`)
  315. }
  316. }
  317. }
  318. /** All mutable lifecycle state for one exact store entry. */
  319. interface SessionEntry {
  320. readonly id: SessionId
  321. readonly session: Session
  322. readonly carrier: Scoped<Session>
  323. readonly emitCtx: Context
  324. announced: boolean
  325. announcing: boolean
  326. appending: boolean
  327. detachRequested: boolean
  328. detach(): void
  329. }
  330. /** Store attachment for the append path; module-private to keep Session store-agnostic publicly. */
  331. const attachments = new WeakMap<Session, SessionEntry>()
  332. /**
  333. * An event-sourced session: an append-only log of {@link SessionEvent}s.
  334. *
  335. * Plain class (not a Service) — create instances via `ctx.sessions.create()`.
  336. * Seeding with an existing event log replays/forks a session.
  337. */
  338. export class Session {
  339. private log: SessionEvent[] = []
  340. /** Single incremental owner of surface acceptance and projection state. */
  341. private readonly surfaceManager = new SurfaceManager(this.log)
  342. /** The ordered surface over this session's event log. */
  343. get surface(): SessionSurface {
  344. return this.surfaceManager
  345. }
  346. /**
  347. * Detached, deep-frozen creation metadata (format version, cwd, lineage,
  348. * seed boundary). Supplied by the store via `ctx.sessions.create()`. When a
  349. * `Session` is constructed bare (tests, ad-hoc replay), a minimal header is
  350. * synthesized (stamped with the current {@link SESSION_FORMAT_VERSION}) so
  351. * `session.header` is always present. Kept out of the event log — it is a
  352. * storage concern, not replayable conversation state.
  353. */
  354. readonly header: SessionHeader
  355. /** The session identity, derived from its durable header's single copy. */
  356. get id(): SessionId {
  357. return this.header.id
  358. }
  359. /**
  360. * The first seq appended IN THIS PROCESS: the length of the constructor
  361. * seed (0 without one). Events below it entered through construction —
  362. * replay, fork, or resume — and were never published on the `session/event`
  363. * firehose (constructor seeds do not emit), so consumers that replay the
  364. * log as a publication substitute (telemetry adoption) start here. Distinct
  365. * from `header.seedLength`, the DURABLE fork-lineage boundary: a resumed
  366. * session's constructor seed is its full stored log, while its header keeps
  367. * the original fork value — this field is the in-process construction fact
  368. * and is deliberately not persisted.
  369. */
  370. readonly firstLiveSeq: number
  371. constructor(id: SessionId, seed?: readonly SessionEvent[], header?: SessionHeader) {
  372. if (seed) {
  373. // Validate the seed to the SAME invariants `append` enforces, so a
  374. // replay/fork (`ctx.sessions.create(id, { seed })`) cannot construct a
  375. // live log that no persistence backend could store: each event's `data`
  376. // must be JSON-serializable, and `seq` must be contiguous from 0 (the
  377. // `seq = log.length` contract the whole system relies on). Without this,
  378. // a bad seed would surface only later as a backend rejection or a silent
  379. // divergence between the live log and disk.
  380. for (const [index, source] of seed.entries()) {
  381. // The seed is a persistence/replay boundary: validate and detach the
  382. // complete event in one lossless-JSON pass.
  383. const snapshot = snapshotJsonValue(source)
  384. if (snapshot === undefined) {
  385. throw new Error(`seed event at index ${index} is not losslessly JSON-serializable`)
  386. }
  387. assertSessionEventEnvelope(snapshot, index)
  388. assertSupportedRequestHeader(snapshot.type, snapshot.data, `seed event at index ${index}`)
  389. if (snapshot.seq !== index) {
  390. throw new Error(`seed event at index ${index} has seq ${snapshot.seq} (expected ${index}); seed must be contiguous from 0`)
  391. }
  392. // A seed is accepted incrementally through the same transition as a
  393. // live append and a full-log fold. The candidate is planned before it
  394. // enters `log`, so a failure cannot partially mutate the surface.
  395. try {
  396. this.surfaceManager.validateNext(snapshot)
  397. } catch (error: unknown) {
  398. throw new Error(`invalid seed event at index ${index}: ${error instanceof Error ? error.message : 'invalid surface metadata'}`)
  399. }
  400. this.log.push(deepFreeze(snapshot))
  401. }
  402. }
  403. this.firstLiveSeq = this.log.length
  404. this.header = snapshotSessionHeader(id, header)
  405. }
  406. /** Cached immutable public snapshot of the private append-only log. */
  407. private eventsSnapshot: readonly SessionEvent[] | undefined
  408. /**
  409. * An immutable snapshot of the append-only event log. The snapshot is reused
  410. * until the next append; a previously returned array does not grow later.
  411. * Events and their nested data are deep-frozen at acceptance, so neither a
  412. * cast nor ordinary JavaScript can rewrite durable history.
  413. */
  414. get events(): readonly SessionEvent[] {
  415. this.eventsSnapshot ??= Object.freeze([...this.log])
  416. return this.eventsSnapshot
  417. }
  418. /** The next event's sequence number — always the log length (the `seq = log.length` contiguity contract). */
  419. get seq(): number {
  420. return this.log.length
  421. }
  422. /**
  423. * Append one typed event to the log and synchronously notify observers via
  424. * the store-owned, module-private publication hooks. The hot path never blocks
  425. * on I/O — persistence plugins buffer asynchronously. Once the event enters
  426. * the log, the append is committed: observer failures are logged and
  427. * contained per listener, so they do not change the return value or prevent
  428. * later listeners from observing the same accepted event.
  429. *
  430. * @param type - The event type (key of {@link SessionEventMap}).
  431. * @param data - The event payload; must be JSON-serializable.
  432. * @param opts - Surface metadata: `surfaceOp` controls how the event enters
  433. * the ordered surface; `sourceEventSeqs` records provenance (the seq
  434. * numbers of events this one derives from). REQUIRED for
  435. * {@link SurfaceEventType} events (every message-producing event must
  436. * declare how it joins the surface, the sole source of derived history) and
  437. * rejected by the compiler for non-surface types like `turn/start` or
  438. * `assistant/chunk`.
  439. * @returns the logged event — its assigned `seq`/`time` plus the SNAPSHOT of
  440. * `data` that entered the log, so reading `event.data` back sees the logged
  441. * value, never the caller's still-mutable input.
  442. * @throws if `data` or surface metadata is not losslessly JSON-serializable
  443. * (BigInt, function, symbol, undefined, negative zero, non-finite number,
  444. * circular reference, sparse array, or an exotic object such as
  445. * Map/Set/Date/class instance), or when the candidate violates the
  446. * canonical surface contract (marker shape and eligibility, unique
  447. * earlier provenance, positional replacement validity, and complete
  448. * shadowed-node coverage). One recursive pass reads, validates, and
  449. * copies each nested value once, so a stateful getter cannot supply one value
  450. * to validation and another to storage. The event log is the durable source
  451. * of truth, so a bad event fails at the append site rather than later during
  452. * a backend flush. A synchronous internal dispatch validation failure or an
  453. * append reentered while this acceptance/publication boundary is open also
  454. * rejects before the log changes.
  455. */
  456. append<T extends SessionEventType>(
  457. type: T,
  458. data: SessionEventMap[T],
  459. ...opts: T extends SurfaceEventType ? [opts: SurfaceIntent] : []
  460. ): SessionEvent<T> {
  461. const surfaceOpts: SurfaceIntent | undefined = opts[0]
  462. const surfaceMetadata = {
  463. ...surfaceOpts?.sourceEventSeqs === undefined ? {} : { sourceEventSeqs: surfaceOpts.sourceEventSeqs },
  464. ...surfaceOpts?.surfaceOp === undefined ? {} : { surfaceOp: surfaceOpts.surfaceOp },
  465. }
  466. const dataSnapshot = snapshotJsonValue(data)
  467. if (dataSnapshot === undefined) {
  468. throw new Error(`session event "${type}" carries non-JSON-serializable data`)
  469. }
  470. assertSupportedRequestHeader(type, dataSnapshot, `session event "${type}"`)
  471. const surfaceMetadataSnapshot = snapshotJsonValue(surfaceMetadata)
  472. if (surfaceMetadataSnapshot === undefined) {
  473. throw new Error(`session event "${type}" carries non-JSON-serializable surface metadata`)
  474. }
  475. const entry = attachments.get(this)
  476. if (entry?.appending) {
  477. throw new Error('session append cannot reenter while another append is being published')
  478. }
  479. const event = deepFreeze({
  480. type,
  481. seq: this.log.length,
  482. time: Date.now(),
  483. data: dataSnapshot,
  484. ...(surfaceMetadataSnapshot as { surfaceOp?: unknown; sourceEventSeqs?: unknown }),
  485. } as unknown as SessionEvent<T>)
  486. this.surfaceManager.validateNext(event as SessionEvent)
  487. if (entry !== undefined) entry.appending = true
  488. try {
  489. let callbacks: SessionCallback[] | undefined
  490. const callbackArgs: unknown[] = [this, event]
  491. if (entry !== undefined) {
  492. callbacks = collectSessionCallbacks(entry.emitCtx, [entry.carrier, 'session/event', ...callbackArgs])
  493. }
  494. this.log.push(event as SessionEvent)
  495. this.eventsSnapshot = undefined
  496. if (callbacks !== undefined && entry !== undefined) {
  497. invokeContainedSessionObservers(entry.emitCtx, 'session/event', entry.id, callbackArgs, callbacks)
  498. }
  499. return event
  500. } finally {
  501. if (entry !== undefined) {
  502. entry.appending = false
  503. if (entry.detachRequested && !entry.announcing) entry.detach()
  504. }
  505. }
  506. }
  507. /** Cached fold of the request-header events — see {@link requestHeader}. */
  508. private headerFold: EpochHeader | undefined
  509. /** Log position (events consumed) the header fold has reached. */
  510. private headerFoldSeq = 0
  511. /**
  512. * The {@link EpochHeader} in force after the log's last header event — the
  513. * header the NEXT request will be compared against — or undefined before
  514. * the first `request/header` snapshot. The live, incrementally-maintained
  515. * form of `foldRequestHeader(session.events)`: each header event is folded
  516. * once, when first seen, so a per-step read costs O(new events).
  517. * @returns the folded header, or undefined when no header event exists yet.
  518. */
  519. requestHeader(): EpochHeader | undefined {
  520. if (this.headerFoldSeq < this.log.length) {
  521. // Frozen on update: the fold is session state exposed by reference — a
  522. // consumer mutating it in place (instead of building a replacement)
  523. // would desync every later comparison against the log, so mutation
  524. // throws instead.
  525. this.headerFold = deepFreeze(foldRequestHeader(this.log.slice(this.headerFoldSeq), this.headerFold))
  526. this.headerFoldSeq = this.log.length
  527. }
  528. return this.headerFold
  529. }
  530. /** The derived-message cache: frozen projections, extended per unseen node. */
  531. private derived: Message[] = []
  532. /** Surface position (nodes projected) the cache has reached. */
  533. private derivedNodes = 0
  534. /** {@link SurfaceManager.replaceGeneration} the cache was built under. */
  535. private derivedGeneration = 0
  536. /**
  537. * Derive the LLM message history by walking the ordered sequences of
  538. * message-producing events maintained by `surfaceOp` markers. The
  539. * surface is the single source of derived history: every message-producing
  540. * append records its `surfaceOp`, so a raw event with no marker (a chunk, a
  541. * turn boundary) is correctly absent, and a compaction `replace` deletes the
  542. * shadowed nodes from the derivation. The projection rules are
  543. * {@link deriveEventMessage}, folded per node.
  544. *
  545. * CACHED: each surface node is projected exactly once, when first seen — a
  546. * call costs O(new nodes), and a surface rewrite (a `replace`;
  547. * {@link SessionSurface.replaceGeneration}) rebuilds. The returned array is
  548. * a fresh snapshot per call (later appends never grow an array a caller
  549. * already holds); the `Message` objects in it are SHARED and **deep-frozen**.
  550. * Their content reuses the already frozen durable event data, so the cache
  551. * needs no second deep clone and consumers still cannot mutate the log.
  552. * @returns a fresh array of the shared, frozen derived history.
  553. */
  554. deriveMessages(): Message[] {
  555. const surface = this.surface
  556. const nodes = surface.nodes
  557. const generation = surface.replaceGeneration
  558. if (generation !== this.derivedGeneration) {
  559. this.derived = []
  560. this.derivedNodes = 0
  561. this.derivedGeneration = generation
  562. }
  563. for (const seq of nodes.slice(this.derivedNodes)) {
  564. // Surface sequences are built from this.log — seq is always a valid
  565. // index by construction. The non-null assertion expresses that invariant.
  566. // eslint-disable-next-line @typescript-eslint/no-non-null-assertion
  567. const msg = this.deriveEventMessage(this.log[seq]!)
  568. // A surface node is one of the five message-producing types, but an
  569. // empty-content assistant/message (a max-tokens step that hosts only
  570. // usage) derives to null and must not enter the transcript.
  571. if (msg) this.derived.push(msg)
  572. }
  573. this.derivedNodes = nodes.length
  574. return [...this.derived]
  575. }
  576. /**
  577. * Project a single event into the LLM message it derives to, or null when
  578. * it produces none — a non-surface event (chunk, boundary, log-only record)
  579. * or an empty-content assistant/message (which exists only to host usage).
  580. * The per-node pure function {@link deriveMessages} folds over the surface;
  581. * an external reconstructor (or the dev invariant) folds the same function
  582. * over a log prefix's surface to rebuild the exact messages any request was
  583. * built from (the reconstructability Agent Note). The returned message is
  584. * the already frozen message nested in the event wrapper and shared by
  585. * delivery, durable history, and model requests.
  586. * @param event - the event to project.
  587. * @returns the derived message, or null when the event produces none.
  588. */
  589. deriveEventMessage(event: SessionEvent): Message | null {
  590. // Intentionally non-exhaustive: only message-producing events derive
  591. // history; turn/step boundaries, chunks, usage, and errors are
  592. // trace/replay data.
  593. switch (event.type) {
  594. // Ordinary prompts, injected context, and mid-turn steering project
  595. // identically in user role: the event's model-facing content stays
  596. // verbatim. Steering's `turn` is log-only. Do NOT
  597. // re-add per-type framing (e.g. `<context>`/`<steering>`) here: framing is
  598. // caller-owned — a producer bakes it into `content`, as workspace-context
  599. // does with `<system-reminder>` — or, if reintroduced, must be driven by
  600. // the event `meta` map and a dedicated renderer, keeping this projection a
  601. // verbatim pass-through. See the deferred design note in
  602. // ../../../../.agents/notes/implemented/simplification/2026-07-20-unwrap-injected-content-envelopes.md
  603. case 'user/message': {
  604. return event.data
  605. }
  606. case 'steering/message': {
  607. return event.data.message
  608. }
  609. case 'assistant/message': {
  610. // Skip an empty-content assistant/message: it exists only to host a
  611. // max-tokens step's usage and must not inject a content-less assistant
  612. // turn into the provider transcript.
  613. if (event.data.message.content.length === 0) return null
  614. return event.data.message
  615. }
  616. case 'tool/result': {
  617. return event.data.message
  618. }
  619. default:
  620. // A non-surface event (boundary, chunk, log-only record) projects to
  621. // no message. Merge-extensible union: no assertNever here.
  622. return null
  623. }
  624. }
  625. }
  626. /** A fork source: either the live session object or its live store id. */
  627. export type SessionForkSource = Session | SessionId
  628. /**
  629. * Rejection codes for session forking: the fork source id is unknown to the
  630. * live store (`SESSION_NOT_FOUND`) or names a session object that is not the
  631. * store's live instance (`SESSION_NOT_LIVE`); the requested child id is
  632. * already taken (`SESSION_ALREADY_EXISTS`); the boundary is not a contiguous
  633. * existing seq (`INVALID_BOUNDARY`); or the selected prefix ends inside an
  634. * open turn (`OPEN_TURN`).
  635. */
  636. export type SessionForkErrorCode =
  637. | 'SESSION_NOT_FOUND'
  638. | 'SESSION_NOT_LIVE'
  639. | 'SESSION_ALREADY_EXISTS'
  640. | 'INVALID_BOUNDARY'
  641. | 'OPEN_TURN'
  642. /** Typed error for session fork rejections. */
  643. export class SessionForkError extends Error {
  644. constructor(message: string, public readonly code: SessionForkErrorCode) {
  645. super(message)
  646. this.name = 'SessionForkError'
  647. }
  648. }
  649. /**
  650. * In-memory session store (`ctx.sessions`).
  651. *
  652. * Persistence is intentionally not implemented here — persistence plugins
  653. * subscribe to `session/event` and flush on `session/flush` / dispose.
  654. */
  655. export class SessionStore extends Service {
  656. private store = new Map<SessionId, SessionEntry>()
  657. private counter = 0
  658. constructor(ctx: Context) {
  659. super(ctx, 'sessions')
  660. }
  661. /**
  662. * Create a session owned by the calling fiber: disposing that fiber stops
  663. * event notification and removes the session from the store. `options.seed`
  664. * populates the session with a copy of those events (replay/fork);
  665. * `options.meta` attaches creation metadata (validated absolute `cwd`, seed
  666. * and parent lineage, and delegation depth) as the immutable
  667. * {@link SessionHeader} (the store fills `version`/`id`/`createdAt`).
  668. *
  669. * For an agent whose session must be torn down IN ORDER with its loop (so the
  670. * loop's final flush is captured before the store attachment ends), do NOT use this
  671. * — fold the session lifecycle into the agent's own effect via
  672. * {@link prepare} + {@link enter} + {@link announce} (see
  673. * `dsh-agent-loop`'s creation transaction).
  674. *
  675. * @param id - the session id; omitted, the store mints `session-<n>`.
  676. * @param options - seed events and/or creation metadata for the header.
  677. * @returns the live session, already entered and announced.
  678. * @throws if a session with `id` already exists, metadata is not a plain
  679. * lossless-JSON record with valid scalar fields, or `meta.cwd` is a
  680. * non-absolute path (storage backends key directories off it).
  681. */
  682. create(id?: SessionId, options?: CreateSessionOptions): Session {
  683. const session = this.prepare(id, options)
  684. // Single effect owned by the calling fiber. Yield the detach BEFORE
  685. // announcing so a throwing `session/created` listener rolls the attach back
  686. // (the generator effect disposes already-yielded disposers on a throw)
  687. // instead of leaking the store entry and its publication hooks.
  688. this.ctx.effect(function* (this: SessionStore) {
  689. yield this.enter(session)
  690. this.announce(session)
  691. }.bind(this), 'sessions.create()')
  692. return session
  693. }
  694. /**
  695. * Build a session WITHOUT entering it into the store — validate the id/cwd and
  696. * construct the {@link Session} (with its immutable {@link SessionHeader}).
  697. * Pairs with {@link enter} + {@link announce}: a caller that owns a composite
  698. * `ctx.effect` (the agent factory) folds the session lifecycle into that ONE
  699. * effect so a fiber unload tears the session + agent down as a single ORDERED
  700. * chain rather than as racing sibling effects — which would remove the publication hooks
  701. * before the loop's closing `session/flush`, dropping the closing events.
  702. *
  703. * @param id - the session id; omitted, the store mints `session-<n>`.
  704. * @param options - seed events and/or creation metadata for the header.
  705. * @returns the constructed session, NOT yet in the store.
  706. * @throws if a session with `id` already exists, metadata is not a plain
  707. * lossless-JSON record with valid scalar fields, or `meta.cwd` is a
  708. * non-absolute path.
  709. */
  710. prepare(id?: SessionId, options?: CreateSessionOptions): Session {
  711. let sessionId: SessionId
  712. if (id === undefined) {
  713. do sessionId = SessionId(`session-${++this.counter}`)
  714. while (this.store.has(sessionId))
  715. } else {
  716. sessionId = SessionId(id)
  717. }
  718. if (this.store.has(sessionId)) throw new Error(`session "${sessionId}" already exists`)
  719. const seed = options?.seed
  720. const meta = options?.meta
  721. const header: SessionHeader = {
  722. version: SESSION_FORMAT_VERSION,
  723. id: sessionId,
  724. createdAt: meta?.createdAt ?? Date.now(),
  725. ...meta?.cwd === undefined ? {} : { cwd: meta.cwd },
  726. ...meta?.parentSession === undefined ? {} : { parentSession: meta.parentSession },
  727. ...meta?.seedLength === undefined ? {} : { seedLength: meta.seedLength },
  728. ...meta?.delegationDepth === undefined ? {} : { delegationDepth: meta.delegationDepth },
  729. }
  730. return new Session(sessionId, seed, header)
  731. }
  732. /**
  733. * Enter a {@link prepare}d session into the store: install the module-private
  734. * append publication hooks and add it to the store. Returns the DETACH
  735. * disposer (hooks + store removal). Does NOT emit `session/created` —
  736. * the caller yields this disposer inside its effect and THEN calls
  737. * {@link announce}, so a throwing `session/created` listener rolls the attach
  738. * back instead of leaking it.
  739. *
  740. * Re-checks the id for a duplicate: `prepare` and `enter` are public
  741. * cross-package primitives and a caller may interleave arbitrary work (or
  742. * another create) between them, so a stale prepared session must NOT overwrite
  743. * a live store entry of the same id — its detach disposer would later delete
  744. * the REAL session. The {@link create} convenience and the agent factory call
  745. * the two back-to-back so they never trip this, but the public seam cannot
  746. * assume that.
  747. *
  748. * @param session - a {@link prepare}d session not yet in the store.
  749. * @returns the detach disposer (publication hooks + store removal). When called from
  750. * a synchronous `session/created` listener, removal and disposal wait until
  751. * that creation dispatch unwinds.
  752. * @throws if a session with this id is already in the store.
  753. */
  754. enter(session: Session): () => void {
  755. const id = session.id
  756. const carrier = scopeTarget(session, scopeOf(this.ctx))
  757. // This is the authoritative collision boundary after arbitrary unpublished
  758. // preparation. Only one exact same-id transaction can publish.
  759. if (this.store.has(id)) throw new Error(`session "${id}" already exists`)
  760. if (attachments.has(session)) throw new Error(`session "${id}" is already attached to a store`)
  761. const entry: SessionEntry = {
  762. id,
  763. session,
  764. carrier,
  765. emitCtx: this.ctx,
  766. announced: false,
  767. announcing: false,
  768. appending: false,
  769. detachRequested: false,
  770. detach: () => { this.detachEntered(entry) },
  771. }
  772. this.store.set(id, entry)
  773. attachments.set(session, entry)
  774. let entered = true
  775. const detach = (): void => {
  776. if (!entered) return
  777. entered = false
  778. // A lifecycle listener may own the advanced detach capability. Keep the
  779. // entry and its publication hooks live until synchronous creation or append
  780. // publication unwinds, then publish the paired disposal edge.
  781. if (entry.announcing || entry.appending) {
  782. entry.detachRequested = true
  783. return
  784. }
  785. entry.detach()
  786. }
  787. return detach
  788. }
  789. /** Remove one exact entered session and emit its paired disposal when announced. */
  790. private detachEntered(entry: SessionEntry): void {
  791. entry.detachRequested = false
  792. // A stale capability cannot remove observers or storage belonging to a
  793. // later same-id lifecycle.
  794. /* v8 ignore next -- enter() rejects replacement while this single-shot detach capability is live. */
  795. if (this.store.get(entry.id) !== entry) return
  796. this.store.delete(entry.id)
  797. attachments.delete(entry.session)
  798. if (entry.announced) this.emitDisposed(entry)
  799. }
  800. /** Emit `session/created` exactly once for an {@link enter}ed session (with
  801. * the carrier {@link enter} captured). Separate from {@link enter} so the
  802. * caller can yield the detach disposer first (rollback safety — see
  803. * {@link enter}).
  804. * @param session - the entered session to announce to listeners.
  805. * @throws if the session is not live or its announcement already began,
  806. * including a reentrant call from a creation listener. */
  807. announce(session: Session): void {
  808. const entry = this.liveEntryFor(session)
  809. if (entry.announced || entry.announcing) {
  810. throw new Error(`session "${entry.id}" was already announced`)
  811. }
  812. // Mark before emit: Cordis emit may deliver to earlier listeners and then
  813. // throw. Rollback must still pair that partial creation with disposal, and
  814. // a listener cannot recursively create a second lifecycle edge.
  815. entry.announced = true
  816. const callbackArgs: unknown[] = [session]
  817. entry.announcing = true
  818. try {
  819. const callbacks = collectSessionCallbacks(this.ctx, [entry.carrier, 'session/created', session])
  820. for (const callback of callbacks) {
  821. // Synchronous throws intentionally propagate and veto publication; the
  822. // yielded detach then emits the paired disposal edge. An async function
  823. // is nevertheless assignable to a void listener, so observe its returned
  824. // promise: rejection is too late to roll back and must be logged instead
  825. // of becoming unhandled.
  826. const returned: unknown = callback(...callbackArgs)
  827. void Promise.resolve(returned).catch((error: unknown) => {
  828. this.ctx.logger.warn(`session "${entry.id}": session/created listener rejected: ${String(error)}`)
  829. })
  830. }
  831. } finally {
  832. entry.announcing = false
  833. if (entry.detachRequested && !entry.appending) entry.detach()
  834. }
  835. }
  836. /** Emit the paired teardown notification with per-listener containment. */
  837. private emitDisposed(entry: SessionEntry): void {
  838. const callbackArgs: unknown[] = [entry.session]
  839. try {
  840. const callbacks = collectSessionCallbacks(this.ctx, [entry.carrier, 'session/disposed', entry.session])
  841. invokeContainedSessionObservers(this.ctx, 'session/disposed', entry.id, callbackArgs, callbacks)
  842. } catch (error: unknown) {
  843. this.ctx.logger.warn(`session "${entry.id}": session/disposed dispatch threw: ${String(error)}`)
  844. }
  845. }
  846. /**
  847. * Dispatch the awaited `session/flush` durability checkpoint for `session`,
  848. * with the carrier captured at {@link enter}. THE flush entry point: the
  849. * store owns the carrier, so callers (the loop's turn-end checkpoint, idle
  850. * injection, teardown drains) must come through here rather than dispatch a
  851. * raw `ctx.parallel('session/flush', …)` — one owner, one spelling, and the
  852. * scoped-dispatch invariant can pin it.
  853. * @param session - the session whose buffered events must reach durable storage.
  854. * @returns resolves when every flush listener has settled; after all settle,
  855. * rejects with the first registered listener failure if any listener failed.
  856. */
  857. async flush(session: Session): Promise<void> {
  858. const { carrier } = this.liveEntryFor(session)
  859. const callbackArgs: unknown[] = [session]
  860. const callbacks = collectSessionCallbacks(this.ctx, [carrier, 'session/flush', session])
  861. const results = await Promise.allSettled(callbacks.map((callback) => {
  862. try {
  863. return callback(...callbackArgs)
  864. } catch (error: unknown) {
  865. // Preserve the listener's exact rejection value; flush is a caller-owned
  866. // failure boundary, and Cordis listeners may throw arbitrary values.
  867. // eslint-disable-next-line @typescript-eslint/prefer-promise-reject-errors
  868. return Promise.reject(error)
  869. }
  870. }))
  871. const failure = results.find((result): result is PromiseRejectedResult => result.status === 'rejected')
  872. if (failure !== undefined) throw failure.reason
  873. }
  874. /** Return the exact live entry; detached/prepared objects reject. */
  875. private liveEntryFor(session: Session): SessionEntry {
  876. const entry = attachments.get(session)
  877. if (entry === undefined || this.store.get(entry.id) !== entry) {
  878. throw new Error(`session "${session.id}" is not live in this store`)
  879. }
  880. return entry
  881. }
  882. /**
  883. * Look up a live session.
  884. * @param id - the session id to look up.
  885. * @returns the session, or undefined when no live session has that id.
  886. */
  887. get(id: SessionId): Session | undefined {
  888. return this.store.get(id)?.session
  889. }
  890. /**
  891. * All live sessions, in creation order.
  892. * @returns a fresh array; mutating it does not affect the store.
  893. */
  894. list(): Session[] {
  895. return [...this.store.values()].map(entry => entry.session)
  896. }
  897. /**
  898. * Create a live child session from a stable prefix of a live source.
  899. * `boundary` is an inclusive source event seq; omitted means the source's
  900. * current last event. The selected slice may end with a between-turn event
  901. * but must not end inside an open turn.
  902. *
  903. * @param source - Live source session object or id.
  904. * @param boundary - Inclusive source event seq to fork through; omitted means
  905. * the source's current last event, and omitted on an empty source forks an
  906. * empty child.
  907. * @param childSessionId - Optional child session id; omitted delegates to
  908. * `SessionStore`'s id policy.
  909. * @returns The created live child session.
  910. */
  911. fork(source: SessionForkSource, boundary?: number, childSessionId?: SessionId): Session {
  912. if (childSessionId !== undefined && this.get(childSessionId) !== undefined) {
  913. throw new SessionForkError(`session "${childSessionId}" already exists`, 'SESSION_ALREADY_EXISTS')
  914. }
  915. const liveSource = this._resolveForkSource(source)
  916. const seed = this._forkSeed(liveSource, boundary)
  917. return this.create(childSessionId, {
  918. seed,
  919. meta: {
  920. ...liveSource.header.cwd !== undefined ? { cwd: liveSource.header.cwd } : {},
  921. parentSession: liveSource.id,
  922. seedLength: seed.length,
  923. },
  924. })
  925. }
  926. private _forkSeed(session: Session, requestedBoundary: number | undefined): SessionEvent[] {
  927. const events = session.events
  928. const lastEvent = events.at(-1)
  929. let boundary: number
  930. if (requestedBoundary !== undefined) {
  931. boundary = requestedBoundary
  932. } else {
  933. if (lastEvent === undefined) return []
  934. boundary = lastEvent.seq
  935. }
  936. if (!Number.isSafeInteger(boundary) || boundary < 0) {
  937. throw new SessionForkError(
  938. `fork boundary for session "${session.id}" must be a non-negative safe integer, got ${String(boundary)}`,
  939. 'INVALID_BOUNDARY',
  940. )
  941. }
  942. if (boundary >= events.length) {
  943. const lastSeq = events.at(-1)?.seq
  944. throw new SessionForkError(
  945. `fork boundary ${boundary} does not exist in session "${session.id}" (last seq: ${lastSeq ?? 'none'})`,
  946. 'INVALID_BOUNDARY',
  947. )
  948. }
  949. const boundaryEvent = events[boundary]
  950. if (boundaryEvent === undefined || boundaryEvent.seq !== boundary) {
  951. throw new SessionForkError(
  952. `fork boundary ${boundary} does not match a contiguous event seq in session "${session.id}"`,
  953. 'INVALID_BOUNDARY',
  954. )
  955. }
  956. const lastTurnBoundary = events.slice(0, boundary + 1)
  957. .findLast(event => event.type === 'turn/start' || event.type === 'turn/end')
  958. if (lastTurnBoundary?.type === 'turn/start') {
  959. throw new SessionForkError(
  960. `fork boundary ${boundary} in session "${session.id}" ends inside open turn ${lastTurnBoundary.data.turn}`,
  961. 'OPEN_TURN',
  962. )
  963. }
  964. return events.slice(0, boundary + 1)
  965. }
  966. private _resolveForkSource(source: SessionForkSource): Session {
  967. if (typeof source === 'string') {
  968. const session = this.get(source)
  969. if (session === undefined) throw new SessionForkError(`session "${source}" not found`, 'SESSION_NOT_FOUND')
  970. return session
  971. }
  972. const live = this.get(source.id)
  973. if (live === undefined) {
  974. throw new SessionForkError(`session "${source.id}" not found`, 'SESSION_NOT_FOUND')
  975. }
  976. if (live !== source) throw new SessionForkError(`session "${source.id}" is not the live store instance`, 'SESSION_NOT_LIVE')
  977. return source
  978. }
  979. }
  980. export default SessionStore