diff --git a/packages/session-projection/session-projection/README.md b/packages/session-projection/session-projection/README.md index af7c9622bb..272c0bf93a 100644 --- a/packages/session-projection/session-projection/README.md +++ b/packages/session-projection/session-projection/README.md @@ -1,33 +1,37 @@ # @deepseek-ai/dsh-session-projection -Session-projection seam. It owns `ctx.sessionProjections`, the registry through which a domain host plugin serves the whole current value of its log-derived per-session state, and through which a carrier (the api-proxy history tail page today; TUI/ACP/headless consumers later) reads every registered value in one synchronous, seq-consistent cut. Design authority: the [session-projection RFC](../../../.agents/notes/proposed/architecture/2026-07-27-session-projection-and-command-log.md). +Session-projection seam. It owns `ctx.sessionProjections`, the registry that DRIVES every registered projection unit forward over committed session events and serves finished whole values to carriers (the api-proxy history tail page and `session/projection` push frame today; TUI/ACP/headless consumers later). A domain registers pure mathematics; the framework owns the drive. Design authority: the [session-projection RFC](../../../.agents/notes/proposed/architecture/2026-07-27-session-projection-and-command-log.md). ## Service: `SessionProjectionRegistry` (ctx key: `sessionProjections`) ### Public API -- `ctx.sessionProjections.register(provider): () => void` Register one domain's provider. Duplicate keys throw; the registration is an effect on the calling fiber, so an unloaded domain plugin's key disappears from subsequent walks (clients read that as capability absence). -- `ctx.sessionProjections.entries(): AnyProjectionProvider[]` Snapshot the registered providers in registration order — the carrier walk surface. +- `ctx.sessionProjections.register(definition): () => void` Register one domain's unit. Duplicate keys and invalid `stateVersion` throw; the registration is an effect on the calling fiber, so an unloaded domain plugin's key (with its cached cells) disappears from subsequent drives and snapshots — clients read that as capability absence. +- `ctx.sessionProjections.onChanged(listener): () => void` Subscribe to the change feed: one call per unit whose state reference changed, per committed event, carrying the schema-validated view and the causing seq. Effect-tied like `register`. +- `ctx.sessionProjections.snapshot(session): ProjectionSnapshot` One consistent synchronous cut over every registered unit — `{ asOfSeq, values }` with `asOfSeq` = the seq of the last event every value reflects (`-1` for an empty log). ### Key Types -- `SessionProjectionMap` — the single merge-extensible type table for the whole chain (host provider, wire block, client cell, React hook). Values are wire-JSON whole values; rendering belongs to the slot system, never this layer. -- `ProjectionProvider` — `{ key, schema, get(agent) }`. `schema` validates the payload before it leaves the host; `get` returns the current whole value and MUST be synchronous. +- `SessionProjectionMap` — the single merge-extensible type table for the whole chain (host unit, wire block, React hook). Values are wire-JSON whole values; rendering belongs to the slot system, never this layer. +- `ProjectionDefinition` — `{ key, schema, init(), apply(state, event), view(state), stateVersion }`: a state-driven computation unit of three pure synchronous functions plus declarations, never an opaque getter. ## Contract -- **Whole-value rule (load-bearing).** A state-carrying log event MUST carry the complete post-change state, never a delta, so the client fold is last-wins by seq. A future domain logging deltas breaks last-wins silently — do not. -- **Synchronous `get`.** Carriers read `session.seq` and every provider value with no await between them; that is what makes `asOfSeq` one consistent cut across all keys. An accidentally-async `get` returns a Promise, which fails the carrier-side `schema.parse` loudly. -- **Full-log view.** `get` runs against the host's full in-memory log (`agent.session.events`); pagination exists only in the history slice served to clients. A last-wins domain may backscan (first hit from the tail terminates); an expensive fold keeps an incremental cache keyed by observed seq. -- **Optional seam.** Domain plugins register under `ctx.inject(['sessionProjections'], …)` so headless assemblies without the registry stay unaffected; carriers use `ctx.get('sessionProjections')` and omit the block entirely when the registry is absent. +- **The framework drives, the domain computes.** The registry subscribes to `session/event` once; every committed event passes every unit's `apply` eagerly. Domains hold no subscriptions. Cells (`{state, observedSeq}` per unit per session, WeakMap-keyed) build lazily — a unit registered after events flowed, or a read of a session predating the registration, folds `init` over the in-memory log on first touch. +- **Same-reference means no work.** `apply` MUST return the same state reference for events that do not concern the unit; the drive gates the change feed on `Object.is`, so non-matching events cost one call and nothing downstream. +- **Whole-value event rule (load-bearing).** A state-carrying log event MUST carry the complete post-change state, never a bare delta — it keeps every transition trivially cheap and every served value self-describing (last-wins for consumers). +- **Synchronous unit discipline.** `init`/`apply`/`view` MUST be synchronous; carriers read `snapshot()` in the same tick as their page slice, which is what makes `asOfSeq` one consistent cut. An accidentally-async `view` returns a Promise, which fails the boundary `schema.parse` loudly. +- **State is plain JSON, `stateVersion` is its invalidation anchor.** The persisted projection cache (a later phase) stores `(sessionId, key, stateVersion, observedSeq, stateJson)` rows; bump `stateVersion` whenever the state shape or the fold semantics change so stale rows are discarded instead of forward-applied into garbage. +- **No wire vocabulary here.** The registry exposes only the change feed and the snapshot read face; carriers (api-proxy) mint their own frames (`session/projection`) and blocks from them. +- **Optional seam.** Domain plugins register under `ctx.inject(['sessionProjections'], …)` so headless assemblies without the registry stay unaffected; carriers use `ctx.get('sessionProjections')` and omit their block/frames entirely when the registry is absent. ## Role -This is the interface package of the capability-seam split: domain host plugins (e.g. `dsh-tool-todo`) contribute providers, carriers (`dsh-host-apiproxy`) consume the walk surface, and neither knows the other. +This is the interface-plus-drive package of the capability-seam split: domain host plugins (e.g. `dsh-tool-todo`) contribute units, carriers (`dsh-host-apiproxy`) consume the snapshot and change feed, and neither knows the other. ## Model Experience -None, as the registry only serves client-facing read models of already-logged session state and touches no prompt, message, schema, stream, or tool result. +None, as the registry only computes client-facing read models of already-logged session state and touches no prompt, message, schema, stream, or tool result. #### KV Cache effect @@ -36,4 +40,6 @@ None; projections never assemble or send provider requests. ## Known Limitations and Deferred Work - **Every tail page carries every registered key** — there is no per-key opt-out or lazy-key request shape yet; acceptable while values are UI-scale whole states (a todo list, a goal snapshot), revisit if a domain's value grows large. -- **Synchronous-`get` discipline is only partially mechanical** — the carrier's `schema.parse` rejects a returned Promise, but a provider that blocks or reads torn non-session state is a review concern; the invariant companion documents why no runtime check exists. +- **Eager drive touches every unit per event** — cheap by construction (whole-value rule, same-reference gate), but a hot path would justify per-unit event-type prefilters, addable without contract change. +- **The persisted projection cache is a later phase** — cells live in memory only; a restart rebuilds by folding the in-memory log on first touch. The `stateVersion` field is the forward-declared invalidation anchor for that phase. +- **Synchronous unit discipline is only partially mechanical** — the boundary `schema.parse` rejects a Promise-returning `view`, but an `apply` that blocks or reads torn non-session state is a review concern; the invariant companion documents why no runtime check exists. diff --git a/packages/session-projection/session-projection/package.json b/packages/session-projection/session-projection/package.json index d4da79599d..473b878f3f 100644 --- a/packages/session-projection/session-projection/package.json +++ b/packages/session-projection/session-projection/package.json @@ -35,13 +35,13 @@ "zod": "^4.4.3" }, "peerDependencies": { - "@deepseek-ai/dsh-agent": "^0.0.1", "@deepseek-ai/dsh-invariants": "^0.0.1", + "@deepseek-ai/dsh-session": "^0.0.1", "cordis": "^4.0.0-rc.7" }, "devDependencies": { - "@deepseek-ai/dsh-agent": "workspace:^", "@deepseek-ai/dsh-invariants": "workspace:^", + "@deepseek-ai/dsh-session": "workspace:^", "cordis": "^4.0.0-rc.7" } } diff --git a/packages/session-projection/session-projection/src/index.ts b/packages/session-projection/session-projection/src/index.ts index 47f66e98ea..8b88c974e8 100644 --- a/packages/session-projection/session-projection/src/index.ts +++ b/packages/session-projection/session-projection/src/index.ts @@ -1,23 +1,25 @@ /** * Session-projection seam: the merge-extensible `SessionProjectionMap` type - * table, the `ProjectionProvider` contract, and the `ctx.sessionProjections` - * registry. Domain host plugins contribute whole current values of - * log-derived per-session state; carriers (api-proxy history tail page, and - * future TUI/ACP consumers) walk the registry synchronously so every key and - * the accompanying `asOfSeq` form one consistent cut. Neither side knows the - * other (capability-seam three-way split). + * table, 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 + * subscription, the per-session watermark cache, and change notification; + * carriers (api-proxy today, TUI/ACP/headless later) consume the snapshot + * read face and the change feed. Neither side knows the other + * (capability-seam three-way split). Design authority: the session-projection + * RFC (.agents/notes/proposed/architecture/2026-07-27-session-projection-and-command-log.md). * - * Whole-value rule (load-bearing): a state-carrying log event MUST carry the - * complete post-change state, never a delta, so the client-side fold is - * last-wins by seq. See the session-projection RFC - * (.agents/notes/proposed/architecture/2026-07-27-session-projection-and-command-log.md). + * Whole-value event rule (load-bearing): a state-carrying log event MUST + * carry the complete post-change state, never a bare delta — it keeps every + * unit's transition trivially cheap and every served value self-describing. * * @module @deepseek-ai/dsh-session-projection */ import { Context, Service } from 'cordis' import type { ZodType } from 'zod' -import type { Agent } from '@deepseek-ai/dsh-agent' +import type { Session, SessionEvent } from '@deepseek-ai/dsh-session' declare module 'cordis' { interface Context { @@ -30,42 +32,110 @@ import type { SessionProjectionMap } from './types.ts' export type { SessionProjectionMap } from './types.ts' /** - * One domain's host-side contribution: the current whole value of its - * log-derived per-session state. + * One domain's state-driven computation unit: three pure synchronous + * functions plus declarations — 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 + * `state` MUST be plain JSON (the persisted-cache precondition). */ -export interface ProjectionProvider { - /** The projection key this provider owns (its `SessionProjectionMap` entry). */ +export interface ProjectionDefinition { + /** The projection key this unit owns (its `SessionProjectionMap` entry). */ key: K - /** Validates the payload before it leaves the host (carriers parse each value through this). */ + /** Validates the wire payload (`view` output) before it leaves the host. */ schema: ZodType /** - * Return the current whole value for one agent's session. MUST be - * synchronous — carriers read `session.seq` and every provider value with no - * await between them, so an async provider would tear the consistency cut - * (an accidentally returned Promise fails the carrier's `schema.parse` - * loudly). Runs against the host's full in-memory log - * (`agent.session.events`): a last-wins domain may backscan from the tail; a - * domain with an expensive fold keeps an incremental cache keyed by observed - * seq. - * @param agent - the agent whose session state is projected. - * @returns the whole current value for this provider's key. + * State for the empty log. + * @returns the initial state. */ - get(agent: Agent): SessionProjectionMap[K] + init(): S + /** + * Pure transition: previous state + one committed event → next state. A + * unit uninterested in an event MUST return the same state reference — an + * unchanged reference (`Object.is`) produces zero downstream work. + * @param state - the state covering all prior events. + * @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] + /** + * Persisted-cache invalidation anchor: bump whenever the state shape or the + * fold semantics change, so persisted `(sessionId, key, stateVersion, + * observedSeq, state)` rows from an older unit are discarded instead of + * being forward-applied into garbage. Non-negative integer. + */ + stateVersion: number } -/** Union-typed view of a registered provider, as seen by carriers walking the table. */ -export type AnyProjectionProvider = ProjectionProvider +/** + * Change-feed listener: one unit's value changed for one session. `value` is + * the schema-validated `view` output; `seq` is the unit's watermark at + * emission (the seq of the event that caused the change). + */ +export type ProjectionChangeListener = ( + session: Session, + key: keyof SessionProjectionMap & string, + value: unknown, + seq: number, +) => void /** - * `ctx.sessionProjections`: the projection provider table. Registration is an - * effect (disposer rides the calling fiber): an unloaded domain plugin's key - * disappears from subsequent walks and clients read it as capability absence. - * Duplicate keys throw. Domain plugins register under - * `ctx.inject(['sessionProjections'], …)` so headless assemblies without the - * registry stay unaffected. + * One consistent read cut over every registered 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. */ + values: Partial +} + +/** Type-erased unit view the drive machinery works with (the register seam already proved the typed contract). */ +interface ErasedDefinition { + key: string + schema: { parse(value: unknown): unknown } + init(): unknown + apply(state: unknown, event: SessionEvent): unknown + view(state: unknown): unknown + stateVersion: number +} + +/** Per-session per-unit watermark cache row. */ +interface UnitCell { + state: unknown + /** Seq of the last event passed through `apply` (regardless of change). */ + observedSeq: number +} + +/** One live registration: the unit plus its per-session cells (dropped whole on disposal). */ +interface Registration { + readonly def: ErasedDefinition + readonly cells: WeakMap +} + +/** + * `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. + * 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 + * calling fiber): an unloaded domain plugin's key disappears from snapshots + * and clients read it as capability absence. Duplicate keys throw. Domain + * plugins register under `ctx.inject(['sessionProjections'], …)` so headless + * assemblies without the registry stay unaffected. */ export class SessionProjectionRegistry extends Service { - private readonly providers = new Map() + private readonly registrations = new Map() + private readonly listeners = new Set() /** * Create and install the registry as `ctx.sessionProjections`. @@ -73,35 +143,107 @@ export class SessionProjectionRegistry extends Service { */ constructor(ctx: Context) { super(ctx, 'sessionProjections') + ctx.on('session/event', (session: Session, event: SessionEvent) => { + this.drive(session, event) + }) } /** - * Register one domain's provider. The registration is an effect on the - * calling context's fiber: disposing the fiber (or calling the returned - * disposer) removes the key from subsequent walks. - * @param provider - key, boundary schema, and synchronous whole-value read. - * @returns the exact disposer that unregisters this provider. + * Register one domain's unit. The registration is an effect on the calling + * context's fiber: disposing the fiber (or calling the returned disposer) + * removes the key — and the unit's cached cells — from subsequent drives + * and snapshots. + * @param definition - key, boundary schema, pure unit functions, and stateVersion. + * @returns the exact disposer that unregisters this unit. */ - register(provider: ProjectionProvider): () => void { + register(definition: ProjectionDefinition): () => void { + 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) { - if (this.providers.has(provider.key)) { - throw new Error(`session projection key ${JSON.stringify(provider.key)} is already registered`) + const key = definition.key as string + if (this.registrations.has(key)) { + throw new Error(`session projection key ${JSON.stringify(key)} is already registered`) } - this.providers.set(provider.key, provider) + this.registrations.set(key, { def: definition as unknown as ErasedDefinition, cells: new WeakMap() }) yield () => { - this.providers.delete(provider.key) + this.registrations.delete(key) } }.bind(this), 'sessionProjections.register()') return () => void dispose() } /** - * Snapshot the registered providers in registration order — the carrier - * walk surface. Each provider carries its own `key` and `schema`. - * @returns the providers registered at this moment. + * 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. + * @returns the exact disposer that unsubscribes. */ - entries(): AnyProjectionProvider[] { - return [...this.providers.values()] + onChanged(listener: ProjectionChangeListener): () => void { + const dispose = this.ctx.effect(() => { + this.listeners.add(listener) + return () => { + this.listeners.delete(listener) + } + }, 'sessionProjections.onChanged()') + return () => void dispose() + } + + /** + * One consistent cut over every registered 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. + * @param session - the session whose projection values are read. + * @returns the snapshot; `values` is empty when no unit is registered. + */ + snapshot(session: Session): ProjectionSnapshot { + const values: Record = {} + for (const registration of this.registrations.values()) { + const cell = this.cellFor(registration, session) + values[registration.def.key] = registration.def.schema.parse(registration.def.view(cell.state)) + } + return { asOfSeq: session.seq - 1, values: values as ProjectionSnapshot['values'] } + } + + /** Fold one unit from init over `events`, producing a cell watermarked at the last folded event. */ + private buildCell(def: ErasedDefinition, events: readonly SessionEvent[]): UnitCell { + let state = def.init() + for (const event of events) state = def.apply(state, event) + return { state, observedSeq: (events.at(-1)?.seq ?? -1) } + } + + /** Read (or lazily build, folding the full in-memory log) one unit's cell. */ + private cellFor(registration: Registration, session: Session): UnitCell { + let cell = registration.cells.get(session) + if (cell === undefined) { + cell = this.buildCell(registration.def, session.events) + registration.cells.set(session, cell) + } + return cell + } + + /** Eager drive: pass one committed event through every registered unit; notify on changed references. */ + private drive(session: Session, event: SessionEvent): void { + for (const registration of this.registrations.values()) { + let cell = registration.cells.get(session) + if (cell === undefined) { + // Late build mid-stream: fold history before this event (seq = log + // index, so the prefix slice is exact), then take the normal gate. + cell = this.buildCell(registration.def, session.events.slice(0, event.seq)) + registration.cells.set(session, cell) + } + const next = registration.def.apply(cell.state, event) + 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)) + for (const listener of this.listeners) { + listener(session, registration.def.key as keyof SessionProjectionMap & string, value, event.seq) + } + } + } } } diff --git a/packages/session-projection/session-projection/src/invariant.ts b/packages/session-projection/session-projection/src/invariant.ts index 36453d72cf..47934c946c 100644 --- a/packages/session-projection/session-projection/src/invariant.ts +++ b/packages/session-projection/session-projection/src/invariant.ts @@ -15,13 +15,16 @@ export const name = 'session-projection-invariant' export const inject = ['invariants'] /** - * No runtime invariant: the registry's own contracts (duplicate-key rejection, - * effect-tied removal) are enforced synchronously at the register() boundary, - * and the served-block relation — every served key has a live registration — - * lives on each carrier's wire path, which emits no cordis event this - * companion could observe; carrier specs assert it instead. Synchronous-`get` - * discipline is enforced as far as practical by the carrier's `schema.parse` - * (a Promise value fails loudly). + * No runtime invariant: the registry's own contracts (duplicate-key and + * stateVersion rejection, effect-tied removal, the Object.is change gate) are + * enforced synchronously inside the service and proven by its spec, the + * drive relation (every committed `session/event` passes every unit) would + * require re-running the drive to check — duplicating the implementation + * rather than detecting drift — and the served-value relation (every served + * key has a live registration) lives on each carrier's wire path, which + * emits no cordis event this companion could observe; carrier specs assert + * it. Synchronous-unit discipline is enforced as far as practical by the + * boundary `schema.parse` (a Promise-returning view fails loudly). */ const install: InvariantInstaller = () => {} diff --git a/packages/session-projection/session-projection/tests/registry.spec.ts b/packages/session-projection/session-projection/tests/registry.spec.ts index d4f193b6cc..bebe17f477 100644 --- a/packages/session-projection/session-projection/tests/registry.spec.ts +++ b/packages/session-projection/session-projection/tests/registry.spec.ts @@ -1,83 +1,190 @@ /** - * SessionProjectionRegistry behavior: registration surfaces through entries(), - * duplicate keys fail loud, and both the returned disposer and the owning - * fiber's disposal remove the key (HMR safety). + * SessionProjectionRegistry unit drive: eager apply on committed events with + * lazy cell build (registration after events, session after registration), + * the Object.is no-change gate (same reference ⇒ zero change-feed work), + * snapshot consistency (asOfSeq = last event seq; values from the watermark + * cache), duplicate-key rejection, stateVersion validation, and effect-tied + * removal of registrations and change listeners (HMR safety). */ import { describe, expect, it } from 'vitest' import { Context } from 'cordis' import { z } from 'zod' -import type { Agent } from '@deepseek-ai/dsh-agent' +import SessionStore from '@deepseek-ai/dsh-session' +import type { Session, SessionEvent } from '@deepseek-ai/dsh-session' import SessionProjectionRegistry from '@deepseek-ai/dsh-session-projection' -import type { ProjectionProvider } from '@deepseek-ai/dsh-session-projection' +import type { ProjectionDefinition } from '@deepseek-ai/dsh-session-projection' -declare module '@deepseek-ai/dsh-session-projection' { +declare module '@deepseek-ai/dsh-session-projection/types' { interface SessionProjectionMap { - 'test/alpha': { value: string } - 'test/beta': number + 'test/marks': { marks: string[] } + 'test/count': number } } -const alphaProvider = (value: string): ProjectionProvider<'test/alpha'> => ({ - key: 'test/alpha', - schema: z.object({ value: z.string() }), - get: () => ({ value }), -}) +declare module '@deepseek-ai/dsh-session' { + interface SessionEventMap { + 'test/mark': { marks: string[] } + } -async function harness(): Promise { - const ctx = new Context() - await ctx.plugin(SessionProjectionRegistry) - return ctx + interface OutOfBandSessionEventMap { + 'test/mark': true + } } -describe('SessionProjectionRegistry', () => { - it('registers a provider, walks it via entries(), and serves get()', async () => { - const ctx = await harness() - ctx.sessionProjections.register(alphaProvider('a')) - const entries = ctx.sessionProjections.entries() - expect(entries.map(entry => entry.key)).toEqual(['test/alpha']) - const provider = entries[0] as ProjectionProvider<'test/alpha'> - expect(provider.get({} as Agent)).toEqual({ value: 'a' }) - expect(provider.schema.parse({ value: 'a' })).toEqual({ value: 'a' }) +/** Whole-value unit: latest test/mark event wins; unrelated events return the same reference. */ +type MarksState = { marks: string[] } | null +const marksUnit = (): ProjectionDefinition<'test/marks', MarksState> => ({ + key: 'test/marks', + schema: z.object({ marks: z.array(z.string()) }), + init: () => null, + apply: (state, event) => (event.type === 'test/mark' ? (event as SessionEvent<'test/mark'>).data : state), + view: state => state ?? { marks: [] }, + stateVersion: 1, +}) + +/** Counting unit over every event — state changes on each apply. */ +const countUnit = (): ProjectionDefinition<'test/count', number> => ({ + key: 'test/count', + schema: z.number().int().nonnegative(), + init: () => 0, + apply: state => state + 1, + view: state => state, + stateVersion: 1, +}) + +async function harness(): Promise<{ ctx: Context; session: Session }> { + const ctx = new Context() + await ctx.plugin(SessionStore) + await ctx.plugin(SessionProjectionRegistry) + return { ctx, session: ctx.sessions.create() } +} + +const mark = (session: Session, marks: string[]): SessionEvent => + session.append('test/mark', { marks }) + +describe('SessionProjectionRegistry drive', () => { + it('drives a registered unit over committed events and snapshots the current value', async () => { + const { ctx, session } = await harness() + ctx.sessionProjections.register(marksUnit()) + mark(session, ['a']) + mark(session, ['a', 'b']) + const snapshot = ctx.sessionProjections.snapshot(session) + expect(snapshot.values['test/marks']).toEqual({ marks: ['a', 'b'] }) + expect(snapshot.asOfSeq).toBe(session.seq - 1) }) - it('preserves registration order across keys', async () => { - const ctx = await harness() - ctx.sessionProjections.register(alphaProvider('a')) - ctx.sessionProjections.register({ - key: 'test/beta', - schema: z.number(), - get: () => 1, + it('builds the cell lazily from the full log for a unit registered after events flowed', async () => { + const { ctx, session } = await harness() + mark(session, ['pre-registration']) + ctx.sessionProjections.register(marksUnit()) + expect(ctx.sessionProjections.snapshot(session).values['test/marks']).toEqual({ marks: ['pre-registration'] }) + // The lazily-built cell then continues on the live drive path. + mark(session, ['after']) + expect(ctx.sessionProjections.snapshot(session).values['test/marks']).toEqual({ marks: ['after'] }) + }) + + it('serves init-derived state and asOfSeq -1 for an empty log', async () => { + const { ctx, session } = await harness() + ctx.sessionProjections.register(marksUnit()) + const snapshot = ctx.sessionProjections.snapshot(session) + expect(snapshot.asOfSeq).toBe(-1) + expect(snapshot.values['test/marks']).toEqual({ marks: [] }) + }) + + it('notifies onChanged with the validated view and the causing seq, and skips same-reference applies', async () => { + const { ctx, session } = await harness() + ctx.sessionProjections.register(marksUnit()) + const seen: { key: string; value: unknown; seq: number; sessionId: string }[] = [] + ctx.sessionProjections.onChanged((changedSession, key, value, seq) => { + seen.push({ key, value, seq, sessionId: String(changedSession.id) }) }) - expect(ctx.sessionProjections.entries().map(entry => entry.key)).toEqual(['test/alpha', 'test/beta']) + const event = mark(session, ['a']) + // Non-matching event: apply returns the same reference — no notification. + session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } }) + expect(seen).toEqual([{ key: 'test/marks', value: { marks: ['a'] }, seq: event.seq, sessionId: String(session.id) }]) }) - it('throws on a duplicate key and keeps the first registration', async () => { - const ctx = await harness() - ctx.sessionProjections.register(alphaProvider('first')) - expect(() => ctx.sessionProjections.register(alphaProvider('second'))) - .toThrow(/"test\/alpha" is already registered/) - const entries = ctx.sessionProjections.entries() - expect(entries).toHaveLength(1) - expect((entries[0] as ProjectionProvider<'test/alpha'>).get({} as Agent)).toEqual({ value: 'first' }) + it('drives independently per session (cells are per-session watermarks)', async () => { + const { ctx, session } = await harness() + const other = ctx.sessions.create() + ctx.sessionProjections.register(marksUnit()) + mark(session, ['one']) + mark(other, ['two']) + expect(ctx.sessionProjections.snapshot(session).values['test/marks']).toEqual({ marks: ['one'] }) + expect(ctx.sessionProjections.snapshot(other).values['test/marks']).toEqual({ marks: ['two'] }) }) - it('register() returns a disposer that removes the key and frees it for re-registration', async () => { - const ctx = await harness() - const dispose = ctx.sessionProjections.register(alphaProvider('a')) + it('runs every registered unit — a changing unit notifies while a same-reference unit stays silent', async () => { + const { ctx, session } = await harness() + ctx.sessionProjections.register(marksUnit()) + ctx.sessionProjections.register(countUnit()) + const changedKeys: string[] = [] + ctx.sessionProjections.onChanged((_session, key) => { + changedKeys.push(key) + }) + session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } }) + // count applied (+1 change), marks returned the same reference. + expect(changedKeys).toEqual(['test/count']) + const snapshot = ctx.sessionProjections.snapshot(session) + expect(snapshot.values['test/count']).toBe(1) + expect(snapshot.values['test/marks']).toEqual({ marks: [] }) + }) + + it('rejects duplicate keys loud and keeps the first unit', async () => { + const { ctx, session } = await harness() + ctx.sessionProjections.register(marksUnit()) + expect(() => ctx.sessionProjections.register(marksUnit())).toThrow(/"test\/marks" is already registered/) + mark(session, ['kept']) + expect(ctx.sessionProjections.snapshot(session).values['test/marks']).toEqual({ marks: ['kept'] }) + }) + + it('rejects a non-integer or negative stateVersion at register time', async () => { + const { ctx } = await harness() + expect(() => ctx.sessionProjections.register({ ...marksUnit(), stateVersion: -1 })).toThrow(/stateVersion/) + expect(() => ctx.sessionProjections.register({ ...marksUnit(), stateVersion: 1.5 })).toThrow(/stateVersion/) + }) + + it('register() disposer removes the key (with its cells) and frees it for re-registration', async () => { + const { ctx, session } = await harness() + const dispose = ctx.sessionProjections.register(marksUnit()) + mark(session, ['cached']) dispose() - expect(ctx.sessionProjections.entries()).toEqual([]) - ctx.sessionProjections.register(alphaProvider('again')) - expect(ctx.sessionProjections.entries()).toHaveLength(1) + expect(ctx.sessionProjections.snapshot(session).values).toEqual({}) + ctx.sessionProjections.register(marksUnit()) + // Fresh registration rebuilds from the log, not from a stale cell. + expect(ctx.sessionProjections.snapshot(session).values['test/marks']).toEqual({ marks: ['cached'] }) }) - it('removes a registration when its owning fiber unloads (HMR safety)', async () => { - const ctx = await harness() + it('removes registrations and change listeners when their owning fiber unloads (HMR safety)', async () => { + const { ctx, session } = await harness() + const notifications: string[] = [] const fiber = await ctx.plugin(Object.assign((inner: Context) => { - inner.sessionProjections.register(alphaProvider('scoped')) + inner.sessionProjections.register(marksUnit()) + inner.sessionProjections.onChanged((_session, key) => { + notifications.push(key) + }) }, { inject: ['sessionProjections'] })) - expect(ctx.sessionProjections.entries()).toHaveLength(1) + mark(session, ['live']) + expect(notifications).toEqual(['test/marks']) await fiber.dispose() - expect(ctx.sessionProjections.entries()).toEqual([]) + mark(session, ['after-dispose']) + expect(notifications).toEqual(['test/marks']) + expect(ctx.sessionProjections.snapshot(session).values).toEqual({}) + }) + + it('fails loud when a unit view violates its own schema (async unit output is unrepresentable)', async () => { + const { ctx, session } = await harness() + ctx.sessionProjections.register({ + key: 'test/marks', + schema: z.object({ marks: z.array(z.string()) }), + init: () => null as MarksState, + apply: state => state, + // A Promise (what an accidentally-async view would return) is not the + // declared shape: the boundary parse rejects it before it leaves. + view: () => Promise.resolve({ marks: [] }) as never, + stateVersion: 1, + }) + expect(() => ctx.sessionProjections.snapshot(session)).toThrow() }) }) diff --git a/packages/session-projection/session-projection/tsconfig.json b/packages/session-projection/session-projection/tsconfig.json index 8b31c9f501..cbd74a19e7 100644 --- a/packages/session-projection/session-projection/tsconfig.json +++ b/packages/session-projection/session-projection/tsconfig.json @@ -15,7 +15,7 @@ "path": "../../../vendor/cordis" }, { - "path": "../../core/agent" + "path": "../../core/session" }, { "path": "../../support/invariants" diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 5d037b6a5b..78d02364d3 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -3338,12 +3338,12 @@ importers: specifier: ^4.4.3 version: 4.4.3 devDependencies: - '@deepseek-ai/dsh-agent': - specifier: workspace:^ - version: link:../../core/agent '@deepseek-ai/dsh-invariants': specifier: workspace:^ version: link:../../support/invariants + '@deepseek-ai/dsh-session': + specifier: workspace:^ + version: link:../../core/session cordis: specifier: ^4.0.0-rc.7 version: 4.0.0-rc.7(@cordisjs/plugin-include@1.0.4)(@cordisjs/plugin-loader@1.0.0-rc.5)