index.ts 52 KB

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