1
/**2
* Service Definition and drive registry for the session-projection capability seam: the merge-extensible state and client-view type3
* tables, the `ProjectionDefinition` state-driven computation unit contract,4
* and the `ctx.sessionProjections` registry that DRIVES every registered unit5
* forward eagerly over committed session events. Domain host plugins6
* contribute pure folds and optional client views; the framework owns the7
* subscription, the per-session watermark cache, and change notification;8
* carriers consume the snapshot read face and the change feed. Neither side9
* knows the other10
* (capability-seam three-way split). Design authority: the session-projection11
* 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 MUST14
* carry the complete post-change state, never a bare delta — it keeps every15
* unit's transition trivially cheap and every served value self-describing.16
*17
* @module @deepseek-ai/dsh-session-projection18
*/20
import { Context, Service } from '@deepseek-ai/cordis'21
import type { ZodType } from 'zod'22
import { SessionLogOffset, SessionSeq } from '@deepseek-ai/dsh-session'23
import type {24
Session,25
SessionEvent,26
SessionHeader,27
SessionSeqCursor,28
} from '@deepseek-ai/dsh-session'30
declare module '@deepseek-ai/cordis' {31
interface Context {32
sessionProjections: SessionProjectionRegistry33
}34
}36
import type { SessionProjectionMap, SessionProjectionStateMap } from './types.ts'38
export type { SessionProjectionMap, SessionProjectionStateMap } from './types.ts'40
/**41
* One domain's state-driven computation unit: a pure synchronous fold plus42
* declarations and an optional client view — never an opaque getter. The framework drives43
* `apply` on every committed session event; the domain holds no44
* subscriptions and owns only the computation. All functions MUST be45
* synchronous (an async unit would tear the carriers' consistency cut), and46
* `state` MUST be plain JSON (the persisted-cache precondition).47
*/48
export 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: K54
/** 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. A65
* unit uninterested in an event MUST return the same state reference — an66
* 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 keeps78
* the two latest raw results and compares them with `Object.is`; an79
* object-valued view must reuse its reference to suppress publication80
* 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
} : never86
/**87
* Persisted-cache invalidation version: bump whenever the serialized state fields or the88
* fold semantics change, so persisted `(sessionId, key, ver, seq, val)`89
* rows from an older unit are discarded instead of being forward-applied90
* into garbage. Non-negative integer.91
*/92
stateVersion: number93
}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 the98
* unit's watermark at emission (the seq of the event that caused the change).99
*/100
export type ProjectionChangeListener = (101
session: Session,102
key: Extract<keyof SessionProjectionMap, string>,103
value: unknown,104
seq: SessionSeq,105
) => void107
/**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 value110
* reflects (`-1` for an empty log).111
*/112
export interface ProjectionSnapshot {113
/** Seq of the last event the values reflect; -1 for an empty log. */114
asOfSeq: SessionSeqCursor115
/** Whole current client value per registered key. */116
values: Partial<SessionProjectionMap>117
}119
/**120
* One unit's checkpoint: its internal state (plain JSON by the unit121
* contract), the seq of the last event folded into it, and the unit122
* `stateVersion` that produced it — the persisted projection-cache row123
* `(sessionId, key, ver, seq, val)` minus the two outer keys. A row is124
* never authoritative, only a fold shortcut: `restore` discards it on a125
* version mismatch or when it claims events past the stored log end.126
*/127
export interface ProjectionCheckpointRow {128
/** The registering unit's `stateVersion` at fold time. */129
ver: number130
/** Seq of the last event folded into `val`; -1 for the empty log. */131
seq: SessionSeqCursor132
/** The unit's internal state — plain JSON per the unit contract. */133
val: unknown134
}136
/** Checkpoint rows keyed by projection key (one session's persisted cache value). */137
export type ProjectionCheckpoint = Record<string, ProjectionCheckpointRow>139
/** Type-erased unit view the drive machinery works with (the registration contract already proved the typed form). */140
interface ErasedDefinition {141
key: string142
stateSchema: { parse(value: unknown): unknown }143
init(header: SessionHeader, inheritedEventCount: SessionLogOffset): unknown144
apply(state: unknown, event: SessionEvent): unknown145
wire: { viewSchema: { parse(value: unknown): unknown }; view(state: unknown): unknown } | undefined146
stateVersion: number147
}149
/** Per-session per-unit watermark and fixed live-drive view buffer. */150
interface UnitCell {151
state: unknown152
/** Seq of the last event passed through `apply` (regardless of change). */153
observedSeq: SessionSeqCursor154
/** `[previousView, currentView]`; undefined slots mean no cached comparison. */155
readonly views: [unknown, unknown]156
}158
/**159
* One live registration: the unit plus its per-session cells (dropped whole160
* once the last registrant releases it).161
*162
* `refs` exists because one unit definition already serves every session — the163
* cells are keyed by `Session` — while registrants are per-session:164
* an agent preset mounts the same tool package once per agent, so N sessions165
* on one preset register the same key N times. Without a count the first166
* registrant would own the disposer, and its session ending would strip the167
* projection from every other live session.168
*/169
interface Registration {170
readonly def: ErasedDefinition171
readonly cells: WeakMap<Session, UnitCell>172
/** Live registrants sharing this unit; the last one out removes the key. */173
refs: number174
}176
/** Convert a log offset to the inclusive cursor immediately before it. */177
function cursorBefore(offset: SessionLogOffset): SessionSeqCursor {178
return offset === 0 ? -1 : SessionSeq(offset - 1)179
}181
/**182
* `ctx.sessionProjections`: the projection unit table and its drive. The183
* service subscribes to `session/event` once; every committed event passes184
* every registered unit's `apply` (eager drive). A changed state reference185
* computes the next client view; the change feed is notified only when its186
* raw result changes by `Object.is`.187
* Cells build lazily — a unit registered after events flowed, or a session188
* older than the registry, folds `init` over the in-memory log on first189
* touch (event or read). Registration is an effect (disposer rides the190
* calling fiber): an unloaded domain plugin's key disappears from snapshots191
* and clients read it as capability absence. A host reader either declares192
* `sessionProjections` in its plugin `inject` or fails explicitly when the193
* registry or required key is absent. Contributors may preserve optional194
* registration through `ctx.inject(['sessionProjections'], ...)`. Registrants sharing a key195
* share one unit and are counted: the same tool package mounted in N agent196
* presets registers N times, and the key survives until the last one197
* unloads.198
*/199
export class SessionProjectionRegistry extends Service {200
private readonly registrations = new Map<string, Registration>()201
private readonly listeners = new Set<ProjectionChangeListener>()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) return211
for (const registration of this.registrations.values()) {212
if (registration.cells.has(session)) continue213
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
}225
/**226
* Register one domain's unit. The registration is an effect on the calling227
* context's fiber: disposing the fiber (or calling the returned disposer)228
* removes the key — and the unit's cached cells — from subsequent drives229
* 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
): () => void241
/**242
* Register one host-only unit. Its state is omitted from client snapshots243
* 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
): () => void253
register<K extends keyof SessionProjectionStateMap, S extends SessionProjectionStateMap[K]>(254
definition: ProjectionDefinition<K, S>,255
): () => void {256
const wire = definition.wire as {257
viewSchema: ZodType258
view(state: S): unknown259
} | undefined260
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 === undefined266
? undefined267
: { 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.key275
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 += 1283
}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) return288
live.refs -= 1289
if (live.refs === 0) this.registrations.delete(key)290
}291
}.bind(this), 'sessionProjections.register()')292
return () => void dispose()293
}295
/**296
* Subscribe to the change feed. The registration is an effect on the297
* 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
}311
/**312
* Read one unit's current host state after materializing every registered313
* 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 undefined325
this.materializeCells(session)326
return this.cellFor(registration, session).state as SessionProjectionStateMap[K]327
}329
/**330
* One consistent cut over every registered client-visible unit for one session, read from331
* the watermark cache (missing cells fold lazily over the in-memory log).332
* Fully synchronous — every value and `asOfSeq` reflect the same log333
* 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) continue347
if (selected !== undefined && !selected.has(registration.def.key)) continue348
const cell = this.cellFor(registration, session)349
values[registration.def.key] = this.viewCell(registration, cell)350
}351
return { asOfSeq: cursorBefore(session.seq), values }352
}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 complete357
* 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 | undefined368
const selected = keys === undefined ? undefined : new Set<string>(keys)369
for (const registration of this.registrations.values()) {370
if (registration.def.wire === undefined) continue371
if (selected !== undefined && !selected.has(registration.def.key)) continue372
const cell = registration.cells.get(session)373
if (cell === undefined) continue374
values[registration.def.key] = this.viewCell(registration, cell)375
if (asOfSeq === undefined || cell.observedSeq < asOfSeq) {376
asOfSeq = cell.observedSeq377
}378
}379
return asOfSeq === undefined ? undefined : { asOfSeq, values }380
}382
/**383
* State-level checkpoint of every persisted unit for one session, read384
* from the watermark cache (missing cells fold lazily over the in-memory385
* log). This is the write side of the persisted projection cache: the386
* returned rows are the `(key → {ver, seq, val})` part of the durable387
* `(sessionId, key, ver, seq, val)`388
* rows. Every `val` is a DETACHED structured clone — never the live389
* cell reference: the watermark cache is this registry's authoritative390
* mutable state, and a caller reaching the live reference could corrupt391
* every subsequent snapshot and frame through it (plain JSON by the unit392
* 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 rows407
}409
/**410
* The stored seq a {@link restore} tail read over `checkpoint` must start411
* at: one event BELOW the lowest usable watermark (a row is usable when412
* its `ver` matches the live unit's `stateVersion`; an absent or mismatched row413
* pulls the floor to `0` — that key must refold the full log). The414
* one-below anchor is load-bearing: the tail then proves how far the415
* stored log still extends, so {@link restore} can detect a log that416
* shrank below a row's watermark (crash-repair truncation) instead of417
* serving the stale row as current — an empty tail read from the anchor418
* yields an end below every watermark and the restore rejects for a full419
* 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 | undefined427
for (const registration of this.registrations.values()) {428
const row = checkpoint[registration.def.key]429
const need = row !== undefined && row.ver === registration.def.stateVersion430
? Math.max(row.seq + 1, 0)431
: 0432
floor = floor === undefined ? need : Math.min(floor, need)433
}434
return floor === undefined ? undefined : SessionLogOffset(Math.max(floor - 1, 0))435
}437
/**438
* View a checkpoint's rows without any log read: for every registered439
* client-visible unit whose row's `ver` matches, serve the schema-validated440
* `view` of the schema-validated stored state; mismatched, malformed, or absent rows leave their key441
* absent (a cold or listing consumer treats it as not-yet-available and a442
* 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.def456
if (def.wire === undefined) continue457
if (selected !== undefined && !selected.has(def.key)) continue458
const row = checkpoint[def.key]459
if (row === undefined || row.ver !== def.stateVersion) continue460
let state: unknown461
try {462
state = def.stateSchema.parse(row.val)463
} catch {464
continue465
}466
values[def.key] = def.wire.viewSchema.parse(def.wire.view(state))467
}468
return values469
}471
/**472
* Cold read: fold every persisted unit over a stored log suffix, seeding473
* each from its checkpoint row when usable — the one read recipe (cached474
* state + forward tail replay + `view`) applied without a live `Session`.475
* Call with the stored events at or past `restoreFloor(checkpoint)` (a476
* `SessionHandle.read` slice) and that same floor as477
* `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 its479
* `ver` matches the live unit's `stateVersion`, it does not predate `baseSeq`480
* (`seq >= baseSeq - 1`), and it does not claim events past the481
* supplied end (`seq <= endSeq`); an unusable row is discarded482
* and its key refolds from `init` — which is only sound over the full483
* log, so a discarded row with `baseSeq > 0` throws (the caller re-reads484
* from seq 0, e.g. after a crash-repair truncation shrank the log below485
* 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 last492
* supplied event's seq, `baseSeq - 1` for an empty tail) plus the493
* 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.def509
const row = checkpoint[def.key]510
const usable = row !== undefined511
&& row.ver === def.stateVersion512
&& row.seq >= beforeBase513
&& row.seq <= endSeq514
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 = usable521
? def.stateSchema.parse(row.val)522
: def.init(header, inheritedEventCount)523
const from = usable ? row.seq : beforeBase524
const startIndex = from - baseSeq + 1525
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
}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 drive545
* 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 = true560
for (const registration of this.registrations.values()) {561
const current = registration.cells.get(session)562
if (current?.observedSeq !== endSeq) {563
complete = false564
break565
}566
}567
if (complete) {568
const values: Record<string, unknown> = {}569
for (const registration of this.registrations.values()) {570
if (registration.def.wire === undefined) continue571
const current = registration.cells.get(session) as UnitCell572
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) continue586
const current = registration.cells.get(session)587
if (current !== undefined && current.observedSeq > row.seq) continue588
registration.cells.set(session, {589
state: row.val,590
observedSeq: row.seq,591
views: [undefined, undefined],592
})593
}594
return restored.snapshot595
}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
}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
}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 cell630
}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) return640
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] = undefined650
}651
cell.state = next652
cell.observedSeq = SessionSeq(seq)653
}654
}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) continue661
if (cell === undefined) {662
// Late build mid-stream: fold history before this event (seq = log663
// 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.state681
const next = registration.def.apply(previousState, event)682
const changed = !Object.is(next, previousState)683
cell.state = next684
cell.observedSeq = event.seq685
const wire = registration.def.wire686
if (changed && wire !== undefined) {687
const views = cell.views688
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] = undefined699
}700
}701
// An unchanged state keeps its current view as the valid comparison702
// value for the next state change.703
}704
}706
/** Return one schema-validated wire value. */707
private viewCell(registration: Registration, cell: UnitCell): unknown {708
const wire = registration.def.wire709
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
}714
export default SessionProjectionRegistry