index.ts 30 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711
  1. /**
  2. * Service Definition and drive registry for the session-projection capability seam: the merge-extensible state and client-view type
  3. * tables, the `ProjectionDefinition` state-driven computation unit contract,
  4. * and the `ctx.sessionProjections` registry that DRIVES every registered unit
  5. * forward eagerly over committed session events. Domain host plugins
  6. * contribute pure folds and optional client views; the framework owns the
  7. * subscription, the per-session watermark cache, and change notification;
  8. * carriers consume the snapshot read face and the change feed. Neither side
  9. * knows the other
  10. * (capability-seam three-way split). Design authority: the session-projection
  11. * RFC (.agents/notes/proposed/architecture/2026-07-27-session-projection-and-command-log.md).
  12. *
  13. * Whole-value event rule (load-bearing): a state-carrying log event MUST
  14. * carry the complete post-change state, never a bare delta — it keeps every
  15. * unit's transition trivially cheap and every served value self-describing.
  16. *
  17. * @module @deepseek-ai/dsh-session-projection
  18. */
  19. import { Context, Service } from '@deepseek-ai/cordis'
  20. import type { ZodType } from 'zod'
  21. import { SessionLogOffset, SessionSeq } from '@deepseek-ai/dsh-session'
  22. import type {
  23. Session,
  24. SessionEvent,
  25. SessionHeader,
  26. SessionSeqCursor,
  27. } from '@deepseek-ai/dsh-session'
  28. declare module '@deepseek-ai/cordis' {
  29. interface Context {
  30. sessionProjections: SessionProjectionRegistry
  31. }
  32. }
  33. import type { SessionProjectionMap, SessionProjectionStateMap } from './types.ts'
  34. export type { SessionProjectionMap, SessionProjectionStateMap } from './types.ts'
  35. /**
  36. * One domain's state-driven computation unit: a pure synchronous fold plus
  37. * declarations and an optional client view — never an opaque getter. The framework drives
  38. * `apply` on every committed session event; the domain holds no
  39. * subscriptions and owns only the computation. All functions MUST be
  40. * synchronous (an async unit would tear the carriers' consistency cut), and
  41. * `state` MUST be plain JSON (the persisted-cache precondition).
  42. */
  43. export interface ProjectionDefinition<
  44. K extends keyof SessionProjectionStateMap,
  45. S extends SessionProjectionStateMap[K] = SessionProjectionStateMap[K],
  46. > {
  47. /** The projection key this unit owns (its `SessionProjectionStateMap` entry). */
  48. key: K
  49. /** Validates persisted state before it seeds a fold. */
  50. stateSchema: ZodType<S>
  51. /**
  52. * State for the empty log and its immutable Session metadata.
  53. * @param header - immutable metadata for the Session being projected.
  54. * @param inheritedEventCount - exact fork-inherited prefix length.
  55. * @returns the initial state.
  56. */
  57. init(header: SessionHeader, inheritedEventCount: SessionLogOffset): NoInfer<S>
  58. /**
  59. * Pure transition: previous state + one committed event → next state. A
  60. * unit uninterested in an event MUST return the same state reference — an
  61. * unchanged reference (`Object.is`) produces zero downstream work.
  62. * @param state - the state covering all prior events.
  63. * @param event - the next committed session event.
  64. * @returns the next state (same reference when the event is not the unit's).
  65. */
  66. apply(state: NoInfer<S>, event: SessionEvent): NoInfer<S>
  67. /** Client view. Omit for host-only units. */
  68. wire?: K extends keyof SessionProjectionMap ? {
  69. /** Validates the wire payload before it leaves the host. */
  70. viewSchema: ZodType<SessionProjectionMap[K]>
  71. /**
  72. * State → wire payload (the read-side projection). The live drive keeps
  73. * the two latest raw results and compares them with `Object.is`; an
  74. * object-valued view must reuse its reference to suppress publication
  75. * across internal-only state changes.
  76. * @param state - the current state.
  77. * @returns the whole current value for this unit's key.
  78. */
  79. view(state: NoInfer<S>): SessionProjectionMap[K]
  80. } : never
  81. /**
  82. * Persisted-cache invalidation version: bump whenever the serialized state fields or the
  83. * fold semantics change, so persisted `(sessionId, key, ver, seq, val)`
  84. * rows from an older unit are discarded instead of being forward-applied
  85. * into garbage. Non-negative integer.
  86. */
  87. stateVersion: number
  88. }
  89. /**
  90. * Change-feed listener: one unit's raw `view` result changed by `Object.is`
  91. * for one session. `value` is the schema-validated output; `seq` is the
  92. * unit's watermark at emission (the seq of the event that caused the change).
  93. */
  94. export type ProjectionChangeListener = (
  95. session: Session,
  96. key: Extract<keyof SessionProjectionMap, string>,
  97. value: unknown,
  98. seq: SessionSeq,
  99. ) => void
  100. /**
  101. * One consistent read cut over every registered client-visible unit for one session.
  102. * `asOfSeq` is the shared watermark — the seq of the last event every value
  103. * reflects (`-1` for an empty log).
  104. */
  105. export interface ProjectionSnapshot {
  106. /** Seq of the last event the values reflect; -1 for an empty log. */
  107. asOfSeq: SessionSeqCursor
  108. /** Whole current client value per registered key. */
  109. values: Partial<SessionProjectionMap>
  110. }
  111. /**
  112. * One unit's checkpoint: its internal state (plain JSON by the unit
  113. * contract), the seq of the last event folded into it, and the unit
  114. * `stateVersion` that produced it — the persisted projection-cache row
  115. * `(sessionId, key, ver, seq, val)` minus the two outer keys. A row is
  116. * never authoritative, only a fold shortcut: `restore` discards it on a
  117. * version mismatch or when it claims events past the stored log end.
  118. */
  119. export interface ProjectionCheckpointRow {
  120. /** The registering unit's `stateVersion` at fold time. */
  121. ver: number
  122. /** Seq of the last event folded into `val`; -1 for the empty log. */
  123. seq: SessionSeqCursor
  124. /** The unit's internal state — plain JSON per the unit contract. */
  125. val: unknown
  126. }
  127. /** Checkpoint rows keyed by projection key (one session's persisted cache value). */
  128. export type ProjectionCheckpoint = Record<string, ProjectionCheckpointRow>
  129. /** Type-erased unit view the drive machinery works with (the registration contract already proved the typed form). */
  130. interface ErasedDefinition {
  131. key: string
  132. stateSchema: { parse(value: unknown): unknown }
  133. init(header: SessionHeader, inheritedEventCount: SessionLogOffset): unknown
  134. apply(state: unknown, event: SessionEvent): unknown
  135. wire: { viewSchema: { parse(value: unknown): unknown }; view(state: unknown): unknown } | undefined
  136. stateVersion: number
  137. }
  138. /** Per-session per-unit watermark and fixed live-drive view buffer. */
  139. interface UnitCell {
  140. state: unknown
  141. /** Seq of the last event passed through `apply` (regardless of change). */
  142. observedSeq: SessionSeqCursor
  143. /** `[previousView, currentView]`; undefined slots mean no cached comparison. */
  144. readonly views: [unknown, unknown]
  145. }
  146. /**
  147. * One live registration: the unit plus its per-session cells (dropped whole
  148. * once the last registrant releases it).
  149. *
  150. * `refs` exists because one unit definition already serves every session — the
  151. * cells are keyed by `Session` — while registrants are per-session:
  152. * an agent preset mounts the same tool package once per agent, so N sessions
  153. * on one preset register the same key N times. Without a count the first
  154. * registrant would own the disposer, and its session ending would strip the
  155. * projection from every other live session.
  156. */
  157. interface Registration {
  158. readonly def: ErasedDefinition
  159. readonly cells: WeakMap<Session, UnitCell>
  160. /** Live registrants sharing this unit; the last one out removes the key. */
  161. refs: number
  162. }
  163. /** Convert a log offset to the inclusive cursor immediately before it. */
  164. function cursorBefore(offset: SessionLogOffset): SessionSeqCursor {
  165. return offset === 0 ? -1 : SessionSeq(offset - 1)
  166. }
  167. /**
  168. * `ctx.sessionProjections`: the projection unit table and its drive. The
  169. * service subscribes to `session/event` once; every committed event passes
  170. * every registered unit's `apply` (eager drive). A changed state reference
  171. * computes the next client view; the change feed is notified only when its
  172. * raw result changes by `Object.is`.
  173. * Cells build lazily — a unit registered after events flowed, or a session
  174. * older than the registry, folds `init` over the in-memory log on first
  175. * touch (event or read). Registration is an effect (disposer rides the
  176. * calling fiber): an unloaded domain plugin's key disappears from snapshots
  177. * and clients read it as capability absence. A host reader either declares
  178. * `sessionProjections` in its plugin `inject` or fails explicitly when the
  179. * registry or required key is absent. Contributors may preserve optional
  180. * registration through `ctx.inject(['sessionProjections'], ...)`. Registrants sharing a key
  181. * share one unit and are counted: the same tool package mounted in N agent
  182. * presets registers N times, and the key survives until the last one
  183. * unloads.
  184. */
  185. export class SessionProjectionRegistry extends Service {
  186. private readonly registrations = new Map<string, Registration>()
  187. private readonly listeners = new Set<ProjectionChangeListener>()
  188. /**
  189. * Create and install the registry as `ctx.sessionProjections`.
  190. * @param ctx - Cordis context that owns the service.
  191. */
  192. constructor(ctx: Context) {
  193. super(ctx, 'sessionProjections')
  194. ctx.on('session/created', (session: Session) => {
  195. if (session.seq !== 0) return
  196. for (const registration of this.registrations.values()) {
  197. if (registration.cells.has(session)) continue
  198. registration.cells.set(session, {
  199. state: registration.def.init(session.header, session.inheritedEventCount),
  200. observedSeq: -1,
  201. views: [undefined, undefined],
  202. })
  203. }
  204. })
  205. ctx.on('session/event', (session: Session, event: SessionEvent) => {
  206. this.drive(session, event)
  207. })
  208. }
  209. /**
  210. * Register one domain's unit. The registration is an effect on the calling
  211. * context's fiber: disposing the fiber (or calling the returned disposer)
  212. * removes the key — and the unit's cached cells — from subsequent drives
  213. * and snapshots.
  214. * @param definition - key, state schema, pure unit functions, and stateVersion.
  215. * @returns the exact disposer that unregisters this unit.
  216. */
  217. register<
  218. K extends keyof SessionProjectionMap,
  219. S extends SessionProjectionStateMap[K],
  220. >(
  221. definition: Omit<ProjectionDefinition<K, S>, 'wire'> & {
  222. wire: NonNullable<ProjectionDefinition<K, S>['wire']>
  223. },
  224. ): () => void
  225. /**
  226. * Register one host-only unit. Its state is omitted from client snapshots
  227. * and always checkpointed like every other unit.
  228. * @param definition - key, state schema, pure unit functions, and stateVersion.
  229. * @returns the exact disposer that unregisters this unit.
  230. */
  231. register<
  232. K extends Exclude<keyof SessionProjectionStateMap, keyof SessionProjectionMap>,
  233. S extends SessionProjectionStateMap[K],
  234. >(
  235. definition: Omit<ProjectionDefinition<K, S>, 'wire'>,
  236. ): () => void
  237. register<K extends keyof SessionProjectionStateMap, S extends SessionProjectionStateMap[K]>(
  238. definition: ProjectionDefinition<K, S>,
  239. ): () => void {
  240. const wire = definition.wire as {
  241. viewSchema: ZodType
  242. view(state: S): unknown
  243. } | undefined
  244. const erased: ErasedDefinition = {
  245. key: definition.key,
  246. stateSchema: definition.stateSchema,
  247. init: (header, inheritedEventCount) => definition.init(header, inheritedEventCount),
  248. apply: (state, event) => definition.apply(state as S, event),
  249. wire: wire === undefined
  250. ? undefined
  251. : { viewSchema: wire.viewSchema, view: state => wire.view(state as S) },
  252. stateVersion: definition.stateVersion,
  253. }
  254. if (!Number.isSafeInteger(definition.stateVersion) || definition.stateVersion < 0) {
  255. throw new Error(`session projection ${JSON.stringify(definition.key)} stateVersion must be a non-negative integer, got ${String(definition.stateVersion)}`)
  256. }
  257. const dispose = this.ctx.effect(function* (this: SessionProjectionRegistry) {
  258. const key = erased.key
  259. const existing = this.registrations.get(key)
  260. if (existing === undefined) {
  261. this.registrations.set(key, { def: erased, cells: new WeakMap(), refs: 1 })
  262. } else {
  263. if (existing.def.stateVersion !== erased.stateVersion) {
  264. throw new Error(`session projection key ${JSON.stringify(key)} is already registered at stateVersion ${String(existing.def.stateVersion)}; refusing to share it with stateVersion ${String(erased.stateVersion)}`)
  265. }
  266. existing.refs += 1
  267. }
  268. yield () => {
  269. const live = this.registrations.get(key)
  270. /* v8 ignore next -- the disposer runs once per successful registration, so the entry it counted is still here */
  271. if (live === undefined) return
  272. live.refs -= 1
  273. if (live.refs === 0) this.registrations.delete(key)
  274. }
  275. }.bind(this), 'sessionProjections.register()')
  276. return () => void dispose()
  277. }
  278. /**
  279. * Subscribe to the change feed. The registration is an effect on the
  280. * calling context's fiber.
  281. * @param listener - called once per client-visible unit whose raw view changed by `Object.is`, per committed event.
  282. * @returns the exact disposer that unsubscribes.
  283. */
  284. onChanged(listener: ProjectionChangeListener): () => void {
  285. const dispose = this.ctx.effect(() => {
  286. this.listeners.add(listener)
  287. return () => {
  288. this.listeners.delete(listener)
  289. }
  290. }, 'sessionProjections.onChanged()')
  291. return () => void dispose()
  292. }
  293. /**
  294. * Read one unit's current host state after materializing every registered
  295. * unit at the Session cursor. Unrelated wire views are not produced.
  296. * The returned value is live; callers must not mutate it.
  297. * @param session - the session whose state is read.
  298. * @param key - the registered unit key.
  299. * @returns current state, or `undefined` when the key is not registered.
  300. */
  301. stateOf<K extends keyof SessionProjectionStateMap>(
  302. session: Session,
  303. key: K,
  304. ): SessionProjectionStateMap[K] | undefined {
  305. const registration = this.registrations.get(key)
  306. if (registration === undefined) return undefined
  307. this.materializeCells(session)
  308. return this.cellFor(registration, session).state as SessionProjectionStateMap[K]
  309. }
  310. /**
  311. * One consistent cut over every registered client-visible unit for one session, read from
  312. * the watermark cache (missing cells fold lazily over the in-memory log).
  313. * Fully synchronous — every value and `asOfSeq` reflect the same log
  314. * position. Each value passes its unit's `viewSchema` before leaving.
  315. * @param session - the session whose projection values are read.
  316. * @param keys - optional client-visible outputs; state materialization remains complete.
  317. * @returns the snapshot; `values` is empty when no selected client-visible unit is registered.
  318. */
  319. snapshot(
  320. session: Session,
  321. keys?: readonly Extract<keyof SessionProjectionMap, string>[],
  322. ): ProjectionSnapshot {
  323. const values: Record<string, unknown> = {}
  324. const selected = keys === undefined ? undefined : new Set<string>(keys)
  325. this.materializeCells(session)
  326. for (const registration of this.registrations.values()) {
  327. if (registration.def.wire === undefined) continue
  328. if (selected !== undefined && !selected.has(registration.def.key)) continue
  329. const cell = this.cellFor(registration, session)
  330. values[registration.def.key] = this.viewCell(registration, cell)
  331. }
  332. return { asOfSeq: cursorBefore(session.seq), values }
  333. }
  334. /**
  335. * Read only already-materialized client-visible cells without folding history.
  336. * Values may trail the live Session and are therefore hints, not a complete
  337. * baseline. Missing cells are omitted.
  338. * @param session - attached Session whose cached cells are inspected.
  339. * @param keys - optional wire keys to view.
  340. * @returns the lowest common cached cut, or `undefined` when no wire cell exists.
  341. */
  342. cachedSnapshot(
  343. session: Session,
  344. keys?: readonly Extract<keyof SessionProjectionMap, string>[],
  345. ): ProjectionSnapshot | undefined {
  346. const values: Record<string, unknown> = {}
  347. let asOfSeq: SessionSeqCursor | undefined
  348. const selected = keys === undefined ? undefined : new Set<string>(keys)
  349. for (const registration of this.registrations.values()) {
  350. if (registration.def.wire === undefined) continue
  351. if (selected !== undefined && !selected.has(registration.def.key)) continue
  352. const cell = registration.cells.get(session)
  353. if (cell === undefined) continue
  354. values[registration.def.key] = this.viewCell(registration, cell)
  355. if (asOfSeq === undefined || cell.observedSeq < asOfSeq) {
  356. asOfSeq = cell.observedSeq
  357. }
  358. }
  359. return asOfSeq === undefined ? undefined : { asOfSeq, values }
  360. }
  361. /**
  362. * State-level checkpoint of every persisted unit for one session, read
  363. * from the watermark cache (missing cells fold lazily over the in-memory
  364. * log). This is the write side of the persisted projection cache: the
  365. * returned rows are the `(key → {ver, seq, val})` part of the durable
  366. * `(sessionId, key, ver, seq, val)`
  367. * rows. Every `val` is a DETACHED structured clone — never the live
  368. * cell reference: the watermark cache is this registry's authoritative
  369. * mutable state, and a caller reaching the live reference could corrupt
  370. * every subsequent snapshot and frame through it (plain JSON by the unit
  371. * contract, so the clone is total).
  372. * @param session - the session whose unit states are checkpointed.
  373. * @returns one row per registered key.
  374. */
  375. checkpoint(session: Session): ProjectionCheckpoint {
  376. const rows: ProjectionCheckpoint = {}
  377. for (const registration of this.registrations.values()) {
  378. const cell = this.cellFor(registration, session)
  379. rows[registration.def.key] = {
  380. ver: registration.def.stateVersion,
  381. seq: cell.observedSeq,
  382. val: structuredClone(cell.state),
  383. }
  384. }
  385. return rows
  386. }
  387. /**
  388. * The stored seq a {@link restore} tail read over `checkpoint` must start
  389. * at: one event BELOW the lowest usable watermark (a row is usable when
  390. * its `ver` matches the live unit's `stateVersion`; an absent or mismatched row
  391. * pulls the floor to `0` — that key must refold the full log). The
  392. * one-below anchor is load-bearing: the tail then proves how far the
  393. * stored log still extends, so {@link restore} can detect a log that
  394. * shrank below a row's watermark (crash-repair truncation) instead of
  395. * serving the stale row as current — an empty tail read from the anchor
  396. * yields an end below every watermark and the restore rejects for a full
  397. * re-read.
  398. * @param checkpoint - persisted rows for one session (possibly stale or empty).
  399. * @returns the offset for the stored-log suffix read (`SessionHandle.read`),
  400. * or `undefined` when no unit is registered (no read needed —
  401. * {@link restore} would serve empty values regardless).
  402. */
  403. restoreFloor(checkpoint: ProjectionCheckpoint): SessionLogOffset | undefined {
  404. let floor: number | undefined
  405. for (const registration of this.registrations.values()) {
  406. const row = checkpoint[registration.def.key]
  407. const need = row !== undefined && row.ver === registration.def.stateVersion
  408. ? Math.max(row.seq + 1, 0)
  409. : 0
  410. floor = floor === undefined ? need : Math.min(floor, need)
  411. }
  412. return floor === undefined ? undefined : SessionLogOffset(Math.max(floor - 1, 0))
  413. }
  414. /**
  415. * View a checkpoint's rows without any log read: for every registered
  416. * client-visible unit whose row's `ver` matches, serve the schema-validated
  417. * `view` of the schema-validated stored state; mismatched, malformed, or absent rows leave their key
  418. * absent (a cold or listing consumer treats it as not-yet-available and a
  419. * fuller read path refolds it). The zero-I/O rung of the read ladder —
  420. * values are as stale as their rows, never wrong.
  421. * @param checkpoint - persisted rows for one session (possibly stale or empty).
  422. * @param keys - optional wire keys to view.
  423. * @returns whole values per key with a usable row; empty when none.
  424. */
  425. viewCheckpoint(
  426. checkpoint: ProjectionCheckpoint,
  427. keys?: readonly Extract<keyof SessionProjectionMap, string>[],
  428. ): Partial<SessionProjectionMap> {
  429. const values: Record<string, unknown> = {}
  430. const selected = keys === undefined ? undefined : new Set<string>(keys)
  431. for (const registration of this.registrations.values()) {
  432. const def = registration.def
  433. if (def.wire === undefined) continue
  434. if (selected !== undefined && !selected.has(def.key)) continue
  435. const row = checkpoint[def.key]
  436. if (row === undefined || row.ver !== def.stateVersion) continue
  437. let state: unknown
  438. try {
  439. state = def.stateSchema.parse(row.val)
  440. } catch {
  441. continue
  442. }
  443. values[def.key] = def.wire.viewSchema.parse(def.wire.view(state))
  444. }
  445. return values
  446. }
  447. /**
  448. * Cold read: fold every persisted unit over a stored log suffix, seeding
  449. * each from its checkpoint row when usable — the one read recipe (cached
  450. * state + forward tail replay + `view`) applied without a live `Session`.
  451. * Call with the stored events at or past `restoreFloor(checkpoint)` (a
  452. * `SessionHandle.read` slice) and that same floor as
  453. * `baseSeq`; the floor's one-below anchor makes the supplied end honest,
  454. * so a shrunk log is detected here. A row is usable iff its
  455. * `ver` matches the live unit's `stateVersion`, it does not predate `baseSeq`
  456. * (`seq >= baseSeq - 1`), and it does not claim events past the
  457. * supplied end (`seq <= endSeq`); an unusable row is discarded
  458. * and its key refolds from `init` — which is only sound over the full
  459. * log, so a discarded row with `baseSeq > 0` throws (the caller re-reads
  460. * from seq 0, e.g. after a crash-repair truncation shrank the log below
  461. * a row's watermark).
  462. * @param checkpoint - persisted rows for one session (possibly stale or empty).
  463. * @param events - the stored events with `seq >= baseSeq`, in seq order.
  464. * @param baseSeq - the seq `events` starts at (its first event's seq when non-empty).
  465. * @param header - immutable metadata for the Session being restored.
  466. * @param inheritedEventCount - exact fork-inherited prefix length supplied to unit initialization.
  467. * @returns the snapshot cut at the supplied log end (`asOfSeq` is the last
  468. * supplied event's seq, `baseSeq - 1` for an empty tail) plus the
  469. * refreshed checkpoint rows at that cut, ready for a durable write-back.
  470. */
  471. restore(
  472. checkpoint: ProjectionCheckpoint,
  473. events: readonly SessionEvent[],
  474. baseSeq: SessionLogOffset,
  475. header: SessionHeader,
  476. inheritedEventCount: SessionLogOffset,
  477. ):
  478. { snapshot: ProjectionSnapshot; checkpoint: ProjectionCheckpoint } {
  479. const endSeq: SessionSeqCursor = events.at(-1)?.seq ?? cursorBefore(baseSeq)
  480. const beforeBase = cursorBefore(baseSeq)
  481. const values: Record<string, unknown> = {}
  482. const refreshed: ProjectionCheckpoint = {}
  483. for (const registration of this.registrations.values()) {
  484. const def = registration.def
  485. const row = checkpoint[def.key]
  486. const usable = row !== undefined
  487. && row.ver === def.stateVersion
  488. && row.seq >= beforeBase
  489. && row.seq <= endSeq
  490. if (!usable && baseSeq > 0) {
  491. throw new Error(
  492. `session projection ${JSON.stringify(def.key)} cannot restore from seq ${baseSeq}: `
  493. + 'its checkpoint row is missing, version-mismatched, or beyond the supplied log end; re-read from seq 0',
  494. )
  495. }
  496. let state = usable
  497. ? def.stateSchema.parse(row.val)
  498. : def.init(header, inheritedEventCount)
  499. const from = usable ? row.seq : beforeBase
  500. const startIndex = from - baseSeq + 1
  501. for (let index = startIndex; index < events.length; index++) {
  502. const event = events[index]
  503. const expectedSeq = SessionSeq(baseSeq + index)
  504. if (event === undefined || event.seq !== expectedSeq) {
  505. throw new Error(`session projection ${JSON.stringify(def.key)} cannot restore across missing seq ${String(expectedSeq)}`)
  506. }
  507. state = def.apply(state, event)
  508. }
  509. if (def.wire !== undefined) values[def.key] = def.wire.viewSchema.parse(def.wire.view(state))
  510. refreshed[def.key] = { ver: def.stateVersion, seq: endSeq, val: state }
  511. }
  512. return {
  513. snapshot: { asOfSeq: endSeq, values: values },
  514. checkpoint: refreshed,
  515. }
  516. }
  517. /**
  518. * Restore an exact cut and install its states on the supplied prepared Session.
  519. * A later publication reuses these cells; ordinary live reads and event drive
  520. * advance any constructor-owned suffix exactly once.
  521. * @param session - exact prepared Session that owns the restored log prefix.
  522. * @param checkpoint - persisted rows for this Session lifecycle.
  523. * @param events - exact events at the observation cut.
  524. * @param baseSeq - first supplied event sequence.
  525. * @returns all projection values at the supplied cut.
  526. */
  527. hydrate(
  528. session: Session,
  529. checkpoint: ProjectionCheckpoint,
  530. events: readonly SessionEvent[],
  531. baseSeq: SessionLogOffset,
  532. ): ProjectionSnapshot {
  533. const endSeq: SessionSeqCursor = events.at(-1)?.seq ?? cursorBefore(baseSeq)
  534. let complete = true
  535. for (const registration of this.registrations.values()) {
  536. const current = registration.cells.get(session)
  537. if (current?.observedSeq !== endSeq) {
  538. complete = false
  539. break
  540. }
  541. }
  542. if (complete) {
  543. const values: Record<string, unknown> = {}
  544. for (const registration of this.registrations.values()) {
  545. if (registration.def.wire === undefined) continue
  546. const current = registration.cells.get(session) as UnitCell
  547. values[registration.def.key] = this.viewCell(registration, current)
  548. }
  549. return { asOfSeq: endSeq, values }
  550. }
  551. const restored = this.restore(
  552. checkpoint,
  553. events,
  554. baseSeq,
  555. session.header,
  556. session.inheritedEventCount,
  557. )
  558. for (const registration of this.registrations.values()) {
  559. const row = restored.checkpoint[registration.def.key]
  560. if (row === undefined) continue
  561. const current = registration.cells.get(session)
  562. if (current !== undefined && current.observedSeq > row.seq) continue
  563. registration.cells.set(session, {
  564. state: row.val,
  565. observedSeq: row.seq,
  566. views: [undefined, undefined],
  567. })
  568. }
  569. return restored.snapshot
  570. }
  571. /** Materialize every registered unit cell at the Session's current cursor. */
  572. private materializeCells(session: Session): void {
  573. for (const registration of this.registrations.values()) this.cellFor(registration, session)
  574. }
  575. /** Fold one unit from init over `events`, producing a cell watermarked at the last folded event. */
  576. private buildCell(
  577. def: ErasedDefinition,
  578. header: SessionHeader,
  579. inheritedEventCount: SessionLogOffset,
  580. events: readonly SessionEvent[],
  581. ): UnitCell {
  582. let state = def.init(header, inheritedEventCount)
  583. for (const event of events) state = def.apply(state, event)
  584. return { state, observedSeq: (events.at(-1)?.seq ?? -1), views: [undefined, undefined] }
  585. }
  586. /** Read (or lazily build, folding the full in-memory log) one unit's cell. */
  587. private cellFor(registration: Registration, session: Session): UnitCell {
  588. let cell = registration.cells.get(session)
  589. if (cell === undefined) {
  590. cell = this.buildCell(
  591. registration.def,
  592. session.header,
  593. session.inheritedEventCount,
  594. session.snapshotEvents(),
  595. )
  596. registration.cells.set(session, cell)
  597. } else {
  598. this.advanceCell(registration.def, cell, session, cursorBefore(session.seq))
  599. }
  600. return cell
  601. }
  602. /** Advance one existing cell through a contiguous Session prefix. */
  603. private advanceCell(
  604. def: ErasedDefinition,
  605. cell: UnitCell,
  606. session: Session,
  607. throughSeq: SessionSeqCursor,
  608. ): void {
  609. if (cell.observedSeq >= throughSeq) return
  610. for (let seq = cell.observedSeq + 1; seq <= throughSeq; seq++) {
  611. const event = session.eventAt(SessionSeq(seq))
  612. if (event === undefined || event.seq !== seq) {
  613. throw new Error(`session projection ${JSON.stringify(def.key)} cannot advance across missing seq ${String(seq)}`)
  614. }
  615. const next = def.apply(cell.state, event)
  616. if (!Object.is(next, cell.state)) {
  617. cell.views[0] = cell.views[1]
  618. cell.views[1] = undefined
  619. }
  620. cell.state = next
  621. cell.observedSeq = SessionSeq(seq)
  622. }
  623. }
  624. /** Eager drive: pass one committed event through every unit; notify on changed raw view references. */
  625. private drive(session: Session, event: SessionEvent): void {
  626. for (const registration of this.registrations.values()) {
  627. let cell = registration.cells.get(session)
  628. if (cell !== undefined && cell.observedSeq >= event.seq) continue
  629. if (cell === undefined) {
  630. // Late build mid-stream: fold history before this event (seq = log
  631. // index, so the prefix slice is exact), then take the normal gate.
  632. cell = this.buildCell(
  633. registration.def,
  634. session.header,
  635. session.inheritedEventCount,
  636. session.snapshotEvents(SessionLogOffset(0), SessionLogOffset(event.seq)),
  637. )
  638. registration.cells.set(session, cell)
  639. } else {
  640. this.advanceCell(
  641. registration.def,
  642. cell,
  643. session,
  644. event.seq === 0 ? -1 : SessionSeq(event.seq - 1),
  645. )
  646. }
  647. const previousState = cell.state
  648. const next = registration.def.apply(previousState, event)
  649. const changed = !Object.is(next, previousState)
  650. cell.state = next
  651. cell.observedSeq = event.seq
  652. const wire = registration.def.wire
  653. if (changed && wire !== undefined) {
  654. const views = cell.views
  655. views[0] = views[1]
  656. if (this.listeners.size > 0) {
  657. views[1] = wire.view(next)
  658. if (!Object.is(views[0], views[1])) {
  659. const value = wire.viewSchema.parse(views[1])
  660. for (const listener of this.listeners) {
  661. listener(session, registration.def.key as Extract<keyof SessionProjectionMap, string>, value, event.seq)
  662. }
  663. }
  664. } else {
  665. views[1] = undefined
  666. }
  667. }
  668. // An unchanged state keeps its current view as the valid comparison
  669. // value for the next state change.
  670. }
  671. }
  672. /** Return one schema-validated wire value. */
  673. private viewCell(registration: Registration, cell: UnitCell): unknown {
  674. const wire = registration.def.wire
  675. if (wire === undefined) throw new Error(`session projection ${JSON.stringify(registration.def.key)} has no wire view`)
  676. return wire.viewSchema.parse(wire.view(cell.state))
  677. }
  678. }
  679. export default SessionProjectionRegistry