1
0

index.ts 55 KB

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