index.ts 41 KB

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