返回源码地图

packages/session/session-projection/src/index.ts

main snapshot · da00f7f5358f · 正文引用章节 06 / 06 / 06;完整原文可核对,不声称全文件人工逐行审计

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

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