|
|
@@ -1,9 +1,9 @@
|
|
|
/**
|
|
|
- * Service Definition and drive registry for the session-projection capability seam: the merge-extensible `SessionProjectionMap` type
|
|
|
- * table, the `ProjectionDefinition` state-driven computation unit contract,
|
|
|
+ * Service Definition and drive registry for the session-projection capability seam: the merge-extensible state and client-view type
|
|
|
+ * tables, the `ProjectionDefinition` state-driven computation unit contract,
|
|
|
* and the `ctx.sessionProjections` registry that DRIVES every registered unit
|
|
|
* forward eagerly over committed session events. Domain host plugins
|
|
|
- * contribute pure mathematics (init/apply/view); the framework owns the
|
|
|
+ * contribute pure folds and optional client views; the framework owns the
|
|
|
* subscription, the per-session watermark cache, and change notification;
|
|
|
* carriers consume the snapshot read face and the change feed. Neither side
|
|
|
* knows the other
|
|
|
@@ -27,28 +27,31 @@ declare module '@deepseek-ai/cordis' {
|
|
|
}
|
|
|
}
|
|
|
|
|
|
-import type { SessionProjectionMap } from './types.ts'
|
|
|
+import type { SessionProjectionMap, SessionProjectionStateMap } from './types.ts'
|
|
|
|
|
|
-export type { SessionProjectionMap } from './types.ts'
|
|
|
+export type { SessionProjectionMap, SessionProjectionStateMap } from './types.ts'
|
|
|
|
|
|
/**
|
|
|
- * One domain's state-driven computation unit: three pure synchronous
|
|
|
- * functions plus declarations — never an opaque getter. The framework drives
|
|
|
+ * One domain's state-driven computation unit: a pure synchronous fold plus
|
|
|
+ * declarations and an optional client view — never an opaque getter. The framework drives
|
|
|
* `apply` on every committed session event; the domain holds no
|
|
|
- * subscriptions and owns only the mathematics. All three functions MUST be
|
|
|
- * synchronous (an async unit would tear the carriers' consistency cut) and
|
|
|
+ * subscriptions and owns only the computation. All functions MUST be
|
|
|
+ * synchronous (an async unit would tear the carriers' consistency cut), and
|
|
|
* `state` MUST be plain JSON (the persisted-cache precondition).
|
|
|
*/
|
|
|
-export interface ProjectionDefinition<K extends keyof SessionProjectionMap, S> {
|
|
|
- /** The projection key this unit owns (its `SessionProjectionMap` entry). */
|
|
|
+export interface ProjectionDefinition<
|
|
|
+ K extends keyof SessionProjectionStateMap,
|
|
|
+ S extends SessionProjectionStateMap[K] = SessionProjectionStateMap[K],
|
|
|
+> {
|
|
|
+ /** The projection key this unit owns (its `SessionProjectionStateMap` entry). */
|
|
|
key: K
|
|
|
- /** Validates the wire payload (`view` output) before it leaves the host. */
|
|
|
- schema: ZodType<SessionProjectionMap[K]>
|
|
|
+ /** Validates persisted state before it seeds a fold. */
|
|
|
+ stateSchema: ZodType<S>
|
|
|
/**
|
|
|
* State for the empty log.
|
|
|
* @returns the initial state.
|
|
|
*/
|
|
|
- init(): S
|
|
|
+ init(): NoInfer<S>
|
|
|
/**
|
|
|
* Pure transition: previous state + one committed event → next state. A
|
|
|
* unit uninterested in an event MUST return the same state reference — an
|
|
|
@@ -57,13 +60,18 @@ export interface ProjectionDefinition<K extends keyof SessionProjectionMap, S> {
|
|
|
* @param event - the next committed session event.
|
|
|
* @returns the next state (same reference when the event is not the unit's).
|
|
|
*/
|
|
|
- apply(state: S, event: SessionEvent): S
|
|
|
- /**
|
|
|
- * State → wire payload (the read-side projection).
|
|
|
- * @param state - the current state.
|
|
|
- * @returns the whole current value for this unit's key.
|
|
|
- */
|
|
|
- view(state: S): SessionProjectionMap[K]
|
|
|
+ apply(state: NoInfer<S>, event: SessionEvent): NoInfer<S>
|
|
|
+ /** Client view. Omit for host-only units. */
|
|
|
+ wire?: K extends keyof SessionProjectionMap ? {
|
|
|
+ /** Validates the wire payload before it leaves the host. */
|
|
|
+ viewSchema: ZodType<SessionProjectionMap[K]>
|
|
|
+ /**
|
|
|
+ * State → wire payload (the read-side projection).
|
|
|
+ * @param state - the current state.
|
|
|
+ * @returns the whole current value for this unit's key.
|
|
|
+ */
|
|
|
+ view(state: NoInfer<S>): SessionProjectionMap[K]
|
|
|
+ } : never
|
|
|
/**
|
|
|
* Persisted-cache invalidation version: bump whenever the serialized state fields or the
|
|
|
* fold semantics change, so persisted `(sessionId, key, ver, seq, val)`
|
|
|
@@ -86,14 +94,14 @@ export type ProjectionChangeListener = (
|
|
|
) => void
|
|
|
|
|
|
/**
|
|
|
- * One consistent read cut over every registered unit for one session.
|
|
|
+ * One consistent read cut over every registered client-visible unit for one session.
|
|
|
* `asOfSeq` is the shared watermark — the seq of the last event every value
|
|
|
* reflects (`-1` for an empty log, mirroring `session/subscribed.lastSeq`).
|
|
|
*/
|
|
|
export interface ProjectionSnapshot {
|
|
|
/** Seq of the last event the values reflect; -1 for an empty log. */
|
|
|
asOfSeq: number
|
|
|
- /** Whole current value per registered key. */
|
|
|
+ /** Whole current client value per registered key. */
|
|
|
values: Partial<SessionProjectionMap>
|
|
|
}
|
|
|
|
|
|
@@ -120,10 +128,10 @@ export type ProjectionCheckpoint = Record<string, ProjectionCheckpointRow>
|
|
|
/** Type-erased unit view the drive machinery works with (the registration contract already proved the typed form). */
|
|
|
interface ErasedDefinition {
|
|
|
key: string
|
|
|
- schema: { parse(value: unknown): unknown }
|
|
|
+ stateSchema: { parse(value: unknown): unknown }
|
|
|
init(): unknown
|
|
|
apply(state: unknown, event: SessionEvent): unknown
|
|
|
- view(state: unknown): unknown
|
|
|
+ wire: { viewSchema: { parse(value: unknown): unknown }; view(state: unknown): unknown } | undefined
|
|
|
stateVersion: number
|
|
|
}
|
|
|
|
|
|
@@ -156,7 +164,8 @@ interface Registration {
|
|
|
* `ctx.sessionProjections`: the projection unit table and its drive. The
|
|
|
* service subscribes to `session/event` once; every committed event passes
|
|
|
* every registered unit's `apply` (eager drive), and a changed state
|
|
|
- * reference notifies the change feed with the schema-validated view.
|
|
|
+ * reference in a client-visible unit notifies the change feed with the
|
|
|
+ * schema-validated view.
|
|
|
* Cells build lazily — a unit registered after events flowed, or a session
|
|
|
* older than the registry, folds `init` over the in-memory log on first
|
|
|
* touch (event or read). Registration is an effect (disposer rides the
|
|
|
@@ -191,22 +200,54 @@ export class SessionProjectionRegistry extends Service {
|
|
|
* @param definition - key, state schema, pure unit functions, and stateVersion.
|
|
|
* @returns the exact disposer that unregisters this unit.
|
|
|
*/
|
|
|
- register<K extends keyof SessionProjectionMap, S>(definition: ProjectionDefinition<K, S>): () => void {
|
|
|
+ register<
|
|
|
+ K extends keyof SessionProjectionMap,
|
|
|
+ S extends SessionProjectionStateMap[K],
|
|
|
+ >(
|
|
|
+ definition: Omit<ProjectionDefinition<K, S>, 'wire'> & {
|
|
|
+ wire: NonNullable<ProjectionDefinition<K, S>['wire']>
|
|
|
+ },
|
|
|
+ ): () => void
|
|
|
+ /**
|
|
|
+ * Register one host-only unit. Its state is omitted from client snapshots
|
|
|
+ * and always checkpointed like every other unit.
|
|
|
+ * @param definition - key, state schema, pure unit functions, and stateVersion.
|
|
|
+ * @returns the exact disposer that unregisters this unit.
|
|
|
+ */
|
|
|
+ register<
|
|
|
+ K extends Exclude<keyof SessionProjectionStateMap, keyof SessionProjectionMap>,
|
|
|
+ S extends SessionProjectionStateMap[K],
|
|
|
+ >(
|
|
|
+ definition: Omit<ProjectionDefinition<K, S>, 'wire'>,
|
|
|
+ ): () => void
|
|
|
+ register<K extends keyof SessionProjectionStateMap, S extends SessionProjectionStateMap[K]>(
|
|
|
+ definition: ProjectionDefinition<K, S>,
|
|
|
+ ): () => void {
|
|
|
+ const wire = definition.wire as {
|
|
|
+ viewSchema: ZodType
|
|
|
+ view(state: S): unknown
|
|
|
+ } | undefined
|
|
|
+ const erased: ErasedDefinition = {
|
|
|
+ key: definition.key,
|
|
|
+ stateSchema: definition.stateSchema,
|
|
|
+ init: () => definition.init(),
|
|
|
+ apply: (state, event) => definition.apply(state as S, event),
|
|
|
+ wire: wire === undefined
|
|
|
+ ? undefined
|
|
|
+ : { viewSchema: wire.viewSchema, view: state => wire.view(state as S) },
|
|
|
+ stateVersion: definition.stateVersion,
|
|
|
+ }
|
|
|
if (!Number.isSafeInteger(definition.stateVersion) || definition.stateVersion < 0) {
|
|
|
throw new Error(`session projection ${JSON.stringify(definition.key)} stateVersion must be a non-negative integer, got ${String(definition.stateVersion)}`)
|
|
|
}
|
|
|
const dispose = this.ctx.effect(function* (this: SessionProjectionRegistry) {
|
|
|
- const key = definition.key as string
|
|
|
+ const key = erased.key
|
|
|
const existing = this.registrations.get(key)
|
|
|
if (existing === undefined) {
|
|
|
- this.registrations.set(key, { def: definition, cells: new WeakMap(), refs: 1 })
|
|
|
+ this.registrations.set(key, { def: erased, cells: new WeakMap(), refs: 1 })
|
|
|
} else {
|
|
|
- // A differing `stateVersion` is the one incompatibility this can name:
|
|
|
- // the versioned contract says the cached state shape differs, so the
|
|
|
- // two registrants cannot share cells. Anything else about a definition
|
|
|
- // is functions, which no runtime comparison can tell apart.
|
|
|
- if (existing.def.stateVersion !== definition.stateVersion) {
|
|
|
- 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(definition.stateVersion)}`)
|
|
|
+ if (existing.def.stateVersion !== erased.stateVersion) {
|
|
|
+ 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)}`)
|
|
|
}
|
|
|
existing.refs += 1
|
|
|
}
|
|
|
@@ -224,7 +265,7 @@ export class SessionProjectionRegistry extends Service {
|
|
|
/**
|
|
|
* Subscribe to the change feed. The registration is an effect on the
|
|
|
* calling context's fiber.
|
|
|
- * @param listener - called once per unit whose state reference changed, per committed event.
|
|
|
+ * @param listener - called once per client-visible unit whose state reference changed, per committed event.
|
|
|
* @returns the exact disposer that unsubscribes.
|
|
|
*/
|
|
|
onChanged(listener: ProjectionChangeListener): () => void {
|
|
|
@@ -238,24 +279,41 @@ export class SessionProjectionRegistry extends Service {
|
|
|
}
|
|
|
|
|
|
/**
|
|
|
- * One consistent cut over every registered unit for one session, read from
|
|
|
+ * Read one unit's current host state without computing unrelated views.
|
|
|
+ * The returned value is live; callers must not mutate it.
|
|
|
+ * @param session - the session whose state is read.
|
|
|
+ * @param key - the registered unit key.
|
|
|
+ * @returns current state, or `undefined` when the key is not registered.
|
|
|
+ */
|
|
|
+ stateOf<K extends keyof SessionProjectionStateMap>(
|
|
|
+ session: Session,
|
|
|
+ key: K,
|
|
|
+ ): SessionProjectionStateMap[K] | undefined {
|
|
|
+ const registration = this.registrations.get(key)
|
|
|
+ if (registration === undefined) return undefined
|
|
|
+ return this.cellFor(registration, session).state as SessionProjectionStateMap[K]
|
|
|
+ }
|
|
|
+
|
|
|
+ /**
|
|
|
+ * One consistent cut over every registered client-visible unit for one session, read from
|
|
|
* the watermark cache (missing cells fold lazily over the in-memory log).
|
|
|
* Fully synchronous — every value and `asOfSeq` reflect the same log
|
|
|
- * position. Each value passes its unit's schema before leaving.
|
|
|
+ * position. Each value passes its unit's `viewSchema` before leaving.
|
|
|
* @param session - the session whose projection values are read.
|
|
|
- * @returns the snapshot; `values` is empty when no unit is registered.
|
|
|
+ * @returns the snapshot; `values` is empty when no client-visible unit is registered.
|
|
|
*/
|
|
|
snapshot(session: Session): ProjectionSnapshot {
|
|
|
const values: Record<string, unknown> = {}
|
|
|
for (const registration of this.registrations.values()) {
|
|
|
+ if (registration.def.wire === undefined) continue
|
|
|
const cell = this.cellFor(registration, session)
|
|
|
- values[registration.def.key] = registration.def.schema.parse(registration.def.view(cell.state))
|
|
|
+ values[registration.def.key] = registration.def.wire.viewSchema.parse(registration.def.wire.view(cell.state))
|
|
|
}
|
|
|
- return { asOfSeq: session.seq - 1, values: values }
|
|
|
+ return { asOfSeq: session.seq - 1, values }
|
|
|
}
|
|
|
|
|
|
/**
|
|
|
- * State-level checkpoint of every registered unit for one session, read
|
|
|
+ * State-level checkpoint of every persisted unit for one session, read
|
|
|
* from the watermark cache (missing cells fold lazily over the in-memory
|
|
|
* log). This is the write side of the persisted projection cache: the
|
|
|
* returned rows are the `(key → {ver, seq, val})` part of the durable
|
|
|
@@ -266,7 +324,7 @@ export class SessionProjectionRegistry extends Service {
|
|
|
* every subsequent snapshot and frame through it (plain JSON by the unit
|
|
|
* contract, so the clone is total).
|
|
|
* @param session - the session whose unit states are checkpointed.
|
|
|
- * @returns one row per registered key; empty when no unit is registered.
|
|
|
+ * @returns one row per registered key.
|
|
|
*/
|
|
|
checkpoint(session: Session): ProjectionCheckpoint {
|
|
|
const rows: ProjectionCheckpoint = {}
|
|
|
@@ -311,8 +369,8 @@ export class SessionProjectionRegistry extends Service {
|
|
|
|
|
|
/**
|
|
|
* View a checkpoint's rows without any log read: for every registered
|
|
|
- * unit whose row's `ver` matches, serve the schema-validated
|
|
|
- * `view` of the stored state; mismatched or absent rows leave their key
|
|
|
+ * client-visible unit whose row's `ver` matches, serve the schema-validated
|
|
|
+ * `view` of the schema-validated stored state; mismatched, malformed, or absent rows leave their key
|
|
|
* absent (a cold or listing consumer treats it as not-yet-available and a
|
|
|
* fuller read path refolds it). The zero-I/O rung of the read ladder —
|
|
|
* values are as stale as their rows, never wrong.
|
|
|
@@ -323,15 +381,22 @@ export class SessionProjectionRegistry extends Service {
|
|
|
const values: Record<string, unknown> = {}
|
|
|
for (const registration of this.registrations.values()) {
|
|
|
const def = registration.def
|
|
|
+ if (def.wire === undefined) continue
|
|
|
const row = checkpoint[def.key]
|
|
|
if (row === undefined || row.ver !== def.stateVersion) continue
|
|
|
- values[def.key] = def.schema.parse(def.view(row.val))
|
|
|
+ let state: unknown
|
|
|
+ try {
|
|
|
+ state = def.stateSchema.parse(row.val)
|
|
|
+ } catch {
|
|
|
+ continue
|
|
|
+ }
|
|
|
+ values[def.key] = def.wire.viewSchema.parse(def.wire.view(state))
|
|
|
}
|
|
|
return values
|
|
|
}
|
|
|
|
|
|
/**
|
|
|
- * Cold read: fold every registered unit over a stored log suffix, seeding
|
|
|
+ * Cold read: fold every persisted unit over a stored log suffix, seeding
|
|
|
* each from its checkpoint row when usable — the one read recipe (cached
|
|
|
* state + forward tail replay + `view`) applied without a live `Session`.
|
|
|
* Call with the events returned by a persistence
|
|
|
@@ -352,7 +417,11 @@ export class SessionProjectionRegistry extends Service {
|
|
|
* supplied event's seq, `baseSeq - 1` for an empty tail) plus the
|
|
|
* refreshed checkpoint rows at that cut, ready for a durable write-back.
|
|
|
*/
|
|
|
- restore(checkpoint: ProjectionCheckpoint, events: readonly SessionEvent[], baseSeq: number):
|
|
|
+ restore(
|
|
|
+ checkpoint: ProjectionCheckpoint,
|
|
|
+ events: readonly SessionEvent[],
|
|
|
+ baseSeq: number,
|
|
|
+ ):
|
|
|
{ snapshot: ProjectionSnapshot; checkpoint: ProjectionCheckpoint } {
|
|
|
const endSeq = events.at(-1)?.seq ?? baseSeq - 1
|
|
|
const values: Record<string, unknown> = {}
|
|
|
@@ -370,12 +439,12 @@ export class SessionProjectionRegistry extends Service {
|
|
|
+ 'its checkpoint row is missing, version-mismatched, or beyond the supplied log end; re-read from seq 0',
|
|
|
)
|
|
|
}
|
|
|
- let state = usable ? row.val : def.init()
|
|
|
+ let state = usable ? def.stateSchema.parse(row.val) : def.init()
|
|
|
const from = usable ? row.seq : baseSeq - 1
|
|
|
for (const event of events) {
|
|
|
if (event.seq > from) state = def.apply(state, event)
|
|
|
}
|
|
|
- values[def.key] = def.schema.parse(def.view(state))
|
|
|
+ if (def.wire !== undefined) values[def.key] = def.wire.viewSchema.parse(def.wire.view(state))
|
|
|
refreshed[def.key] = { ver: def.stateVersion, seq: endSeq, val: state }
|
|
|
}
|
|
|
return {
|
|
|
@@ -415,8 +484,8 @@ export class SessionProjectionRegistry extends Service {
|
|
|
const changed = !Object.is(next, cell.state)
|
|
|
cell.state = next
|
|
|
cell.observedSeq = event.seq
|
|
|
- if (changed && this.listeners.size > 0) {
|
|
|
- const value = registration.def.schema.parse(registration.def.view(next))
|
|
|
+ if (changed && registration.def.wire !== undefined && this.listeners.size > 0) {
|
|
|
+ const value = registration.def.wire.viewSchema.parse(registration.def.wire.view(next))
|
|
|
for (const listener of this.listeners) {
|
|
|
listener(session, registration.def.key as Extract<keyof SessionProjectionMap, string>, value, event.seq)
|
|
|
}
|