|
|
@@ -1,13 +1,13 @@
|
|
|
/**
|
|
|
* 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).
|
|
|
+ * the Object.is no-change gates (same state or raw view 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 { describe, expect, it, vi } from 'vitest'
|
|
|
import { Context } from '@deepseek-ai/cordis'
|
|
|
import { z } from 'zod'
|
|
|
import SessionStore, { SessionId } from '@deepseek-ai/dsh-session'
|
|
|
@@ -19,14 +19,12 @@ declare module '@deepseek-ai/dsh-session-projection/types' {
|
|
|
interface SessionProjectionStateMap {
|
|
|
'test/marks': MarksState
|
|
|
'test/count': number
|
|
|
- 'test/buffered': { marks: string[]; draft: string }
|
|
|
- 'test/label': string
|
|
|
+ 'test/stable-view': StableViewState
|
|
|
}
|
|
|
|
|
|
interface SessionProjectionMap {
|
|
|
'test/marks': { marks: string[] }
|
|
|
- 'test/buffered': string[]
|
|
|
- 'test/label': string
|
|
|
+ 'test/stable-view': { marks: string[] }
|
|
|
}
|
|
|
}
|
|
|
|
|
|
@@ -36,7 +34,15 @@ declare module '@deepseek-ai/dsh-session/types' {
|
|
|
}
|
|
|
}
|
|
|
|
|
|
-type MarksState = { marks: string[] } | null
|
|
|
+interface MarksView {
|
|
|
+ marks: string[]
|
|
|
+}
|
|
|
+type MarksState = MarksView | null
|
|
|
+interface StableViewState {
|
|
|
+ revision: number
|
|
|
+ value: MarksView
|
|
|
+}
|
|
|
+const marksViewSchema: z.ZodType<MarksView> = z.object({ marks: z.array(z.string()) })
|
|
|
const RESTORE_HEADER: SessionHeader = {
|
|
|
version: 0,
|
|
|
id: SessionId('projection-restore'),
|
|
|
@@ -46,11 +52,11 @@ const RESTORE_HEADER: SessionHeader = {
|
|
|
const marksUnit = (): Omit<ProjectionDefinition<'test/marks', MarksState>, 'wire'>
|
|
|
& { wire: NonNullable<ProjectionDefinition<'test/marks', MarksState>['wire']> } => ({
|
|
|
key: 'test/marks',
|
|
|
- stateSchema: z.object({ marks: z.array(z.string()) }).nullable(),
|
|
|
+ stateSchema: marksViewSchema.nullable(),
|
|
|
init: () => null,
|
|
|
apply: (state, event) => (event.type === 'test/mark' ? (event).data : state),
|
|
|
wire: {
|
|
|
- viewSchema: z.object({ marks: z.array(z.string()) }),
|
|
|
+ viewSchema: marksViewSchema,
|
|
|
view: state => state ?? { marks: [] },
|
|
|
},
|
|
|
stateVersion: 1,
|
|
|
@@ -65,6 +71,27 @@ const countUnit = (): ProjectionDefinition<'test/count', number> => ({
|
|
|
stateVersion: 1,
|
|
|
})
|
|
|
|
|
|
+const stableViewUnit = (
|
|
|
+ view: (state: StableViewState) => StableViewState['value'],
|
|
|
+) => ({
|
|
|
+ key: 'test/stable-view',
|
|
|
+ stateSchema: z.object({
|
|
|
+ revision: z.number().int().nonnegative(),
|
|
|
+ value: marksViewSchema,
|
|
|
+ }),
|
|
|
+ init: () => ({ revision: 0, value: { marks: [] } }),
|
|
|
+ apply: (state, event) => {
|
|
|
+ if (event.type === 'turn/start') return { ...state, revision: state.revision + 1 }
|
|
|
+ if (event.type === 'test/mark') return { revision: state.revision + 1, value: event.data }
|
|
|
+ return state
|
|
|
+ },
|
|
|
+ wire: {
|
|
|
+ viewSchema: marksViewSchema,
|
|
|
+ view,
|
|
|
+ },
|
|
|
+ stateVersion: 1,
|
|
|
+}) satisfies ProjectionDefinition<'test/stable-view', StableViewState>
|
|
|
+
|
|
|
async function harness(): Promise<{ ctx: Context; session: Session }> {
|
|
|
const ctx = new Context()
|
|
|
await ctx.plugin(SessionStore)
|
|
|
@@ -117,111 +144,62 @@ describe('SessionProjectionRegistry drive', () => {
|
|
|
expect(seen).toEqual([{ key: 'test/marks', value: { marks: ['a'] }, seq: event.seq, sessionId: String(session.id) }])
|
|
|
})
|
|
|
|
|
|
- it('keeps the feed quiet while a changed state serves an identity-stable view (draft buffering)', async () => {
|
|
|
+ it('does not compute a view while no change listener exists', async () => {
|
|
|
const { ctx, session } = await harness()
|
|
|
- // Unit buffering a working field beside its wire array: the view projects
|
|
|
- // only `marks`, whose identity survives draft-only applies.
|
|
|
- ctx.sessionProjections.register({
|
|
|
- key: 'test/buffered',
|
|
|
- stateSchema: z.object({ marks: z.array(z.string()), draft: z.string() }),
|
|
|
- init: () => ({ marks: [], draft: '' }),
|
|
|
- apply: (state, event) => {
|
|
|
- if (event.type === 'test/mark') return { marks: event.data.marks, draft: '' }
|
|
|
- if (event.type === 'turn/start') return { marks: state.marks, draft: `draft-${String(event.seq)}` }
|
|
|
- return state
|
|
|
- },
|
|
|
- wire: { viewSchema: z.array(z.string()), view: state => state.marks },
|
|
|
- stateVersion: 1,
|
|
|
- })
|
|
|
- const seen: { value: unknown; seq: number }[] = []
|
|
|
- ctx.sessionProjections.onChanged((_session, key, value, seq) => {
|
|
|
- if (key === 'test/buffered') seen.push({ value, seq })
|
|
|
- })
|
|
|
- const first = mark(session, ['a'])
|
|
|
- // Draft-only applies change the state reference but not the served view.
|
|
|
+ const view = vi.fn((state: StableViewState) => state.value)
|
|
|
+ ctx.sessionProjections.register(stableViewUnit(view))
|
|
|
+
|
|
|
session.append('turn/start', { turn: 1 })
|
|
|
session.append('turn/start', { turn: 2 })
|
|
|
- const second = mark(session, ['a', 'b'])
|
|
|
- expect(seen).toEqual([
|
|
|
- { value: ['a'], seq: first.seq },
|
|
|
- { value: ['a', 'b'], seq: second.seq },
|
|
|
- ])
|
|
|
- // The quiet applies still advanced the state itself.
|
|
|
- expect(ctx.sessionProjections.stateOf(session, 'test/buffered')?.draft).toBe('')
|
|
|
- expect(ctx.sessionProjections.snapshot(session).values['test/buffered']).toEqual(['a', 'b'])
|
|
|
+
|
|
|
+ expect(ctx.sessionProjections.stateOf(session, 'test/stable-view')?.revision).toBe(2)
|
|
|
+ expect(view).not.toHaveBeenCalled()
|
|
|
})
|
|
|
|
|
|
- it("computes each distinct state's view once: the memo serves previous states to the gate and snapshots", async () => {
|
|
|
+ it('publishes the first observed view and suppresses later same-reference views', async () => {
|
|
|
const { ctx, session } = await harness()
|
|
|
- let viewCalls = 0
|
|
|
- ctx.sessionProjections.register({
|
|
|
- key: 'test/buffered',
|
|
|
- stateSchema: z.object({ marks: z.array(z.string()), draft: z.string() }),
|
|
|
- init: () => ({ marks: [], draft: '' }),
|
|
|
- apply: (state, event) => {
|
|
|
- if (event.type === 'test/mark') return { marks: event.data.marks, draft: '' }
|
|
|
- if (event.type === 'turn/start') return { marks: state.marks, draft: `draft-${String(event.seq)}` }
|
|
|
- return state
|
|
|
- },
|
|
|
- wire: {
|
|
|
- viewSchema: z.array(z.string()),
|
|
|
- view: (state) => {
|
|
|
- viewCalls += 1
|
|
|
- return state.marks
|
|
|
- },
|
|
|
- },
|
|
|
- stateVersion: 1,
|
|
|
- })
|
|
|
+ const view = vi.fn((state: StableViewState) => state.value)
|
|
|
+ ctx.sessionProjections.register(stableViewUnit(view))
|
|
|
+
|
|
|
const seen: unknown[] = []
|
|
|
ctx.sessionProjections.onChanged((_session, key, value) => {
|
|
|
- if (key === 'test/buffered') seen.push(value)
|
|
|
+ if (key === 'test/stable-view') seen.push(value)
|
|
|
})
|
|
|
- // First change touches two never-seen states (init and next): two calls.
|
|
|
- mark(session, ['a'])
|
|
|
- expect(viewCalls).toBe(2)
|
|
|
- // Draft-only change: the new state computes, the previous is a memo hit.
|
|
|
+
|
|
|
session.append('turn/start', { turn: 1 })
|
|
|
- expect(viewCalls).toBe(3)
|
|
|
- mark(session, ['a', 'b'])
|
|
|
- expect(viewCalls).toBe(4)
|
|
|
- expect(seen).toEqual([['a'], ['a', 'b']])
|
|
|
- // Snapshot reads reuse the same memo instead of recomputing the view.
|
|
|
- expect(ctx.sessionProjections.snapshot(session).values['test/buffered']).toEqual(['a', 'b'])
|
|
|
- expect(viewCalls).toBe(4)
|
|
|
+ session.append('turn/start', { turn: 2 })
|
|
|
+
|
|
|
+ expect(seen).toEqual([{ marks: [] }])
|
|
|
+ expect(view).toHaveBeenCalledTimes(2)
|
|
|
+
|
|
|
+ mark(session, ['changed'])
|
|
|
+ expect(seen).toEqual([{ marks: [] }, { marks: ['changed'] }])
|
|
|
+ expect(view).toHaveBeenCalledTimes(3)
|
|
|
})
|
|
|
|
|
|
- it('keeps dedup honest across listener generations: a return to an old value after an unobserved change still fires', async () => {
|
|
|
+ it('publishes the first view after an unobserved state change', async () => {
|
|
|
const { ctx, session } = await harness()
|
|
|
- // A primitive-valued view compares by value under Object.is (the title
|
|
|
- // unit's shape), which is exactly where remembering a delivered value —
|
|
|
- // instead of comparing the two states in hand — would silence a real
|
|
|
- // transition.
|
|
|
- ctx.sessionProjections.register({
|
|
|
- key: 'test/label',
|
|
|
- stateSchema: z.string(),
|
|
|
- init: () => '',
|
|
|
- apply: (state, event) => (event.type === 'test/mark' ? event.data.marks[0] ?? '' : state),
|
|
|
- wire: { viewSchema: z.string(), view: state => state },
|
|
|
- stateVersion: 1,
|
|
|
- })
|
|
|
- const first: string[] = []
|
|
|
+ const view = vi.fn((state: StableViewState) => state.value)
|
|
|
+ ctx.sessionProjections.register(stableViewUnit(view))
|
|
|
+ const first: unknown[] = []
|
|
|
const stop = ctx.sessionProjections.onChanged((_session, key, value) => {
|
|
|
- if (key === 'test/label') first.push(value as string)
|
|
|
+ if (key === 'test/stable-view') first.push(value)
|
|
|
})
|
|
|
- mark(session, ['A'])
|
|
|
+
|
|
|
+ session.append('turn/start', { turn: 1 })
|
|
|
stop()
|
|
|
- // Unobserved transition away from 'A'…
|
|
|
- mark(session, ['B'])
|
|
|
- // …then a new listener generation subscribes and the value returns:
|
|
|
- // dedup memory frozen at the delivered 'A' would silence this delivery;
|
|
|
- // the per-step previous-state comparison sees 'B' → 'A' and fires.
|
|
|
- const second: string[] = []
|
|
|
+ session.append('turn/start', { turn: 2 })
|
|
|
+ expect(view).toHaveBeenCalledTimes(1)
|
|
|
+
|
|
|
+ const resumed: unknown[] = []
|
|
|
ctx.sessionProjections.onChanged((_session, key, value) => {
|
|
|
- if (key === 'test/label') second.push(value as string)
|
|
|
+ if (key === 'test/stable-view') resumed.push(value)
|
|
|
})
|
|
|
- mark(session, ['A'])
|
|
|
- expect(first).toEqual(['A'])
|
|
|
- expect(second).toEqual(['A'])
|
|
|
+ session.append('turn/start', { turn: 3 })
|
|
|
+
|
|
|
+ expect(first).toEqual([{ marks: [] }])
|
|
|
+ expect(resumed).toEqual([{ marks: [] }])
|
|
|
+ expect(view).toHaveBeenCalledTimes(2)
|
|
|
})
|
|
|
|
|
|
it('drives independently per session (cells are per-session watermarks)', async () => {
|