|
|
@@ -34,7 +34,7 @@ declare module '@deepseek-ai/dsh-session/types' {
|
|
|
|
|
|
type MarksState = { marks: string[] } | null
|
|
|
/** Whole-value unit: latest test/mark event wins; unrelated events return the same reference. */
|
|
|
-const marksUnit = (): ProjectionDefinition<'test/marks', MarksState>
|
|
|
+const marksUnit = (): Omit<ProjectionDefinition<'test/marks', MarksState>, 'wire' | 'persist'>
|
|
|
& { wire: NonNullable<ProjectionDefinition<'test/marks', MarksState>['wire']> } => ({
|
|
|
key: 'test/marks',
|
|
|
stateSchema: z.object({ marks: z.array(z.string()) }).nullable(),
|
|
|
@@ -109,90 +109,6 @@ describe('SessionProjectionRegistry drive', () => {
|
|
|
expect(seen).toEqual([{ key: 'test/marks', value: { marks: ['a'] }, seq: event.seq, sessionId: String(session.id) }])
|
|
|
})
|
|
|
|
|
|
- it('contains a throwing change listener and continues the remaining feed', async () => {
|
|
|
- const { ctx, session } = await harness()
|
|
|
- ctx.sessionProjections.register(marksUnit())
|
|
|
- const seen: string[] = []
|
|
|
- ctx.sessionProjections.onChanged(() => {
|
|
|
- throw new Error('change listener boom')
|
|
|
- })
|
|
|
- ctx.sessionProjections.onChanged((_session, key) => {
|
|
|
- seen.push(key)
|
|
|
- })
|
|
|
-
|
|
|
- expect(() => mark(session, ['contained'])).not.toThrow()
|
|
|
- expect(seen).toEqual(['test/marks'])
|
|
|
- })
|
|
|
-
|
|
|
- it('publishes only client-visible changes and honors the direct disposer', async () => {
|
|
|
- const { ctx, session } = await harness()
|
|
|
- ctx.sessionProjections.register(marksUnit())
|
|
|
- ctx.sessionProjections.register(countUnit())
|
|
|
- const seen: string[] = []
|
|
|
- const dispose = ctx.sessionProjections.onChanged((_session, key) => {
|
|
|
- seen.push(key)
|
|
|
- })
|
|
|
- mark(session, ['wire'])
|
|
|
- expect(seen).toEqual(['test/marks'])
|
|
|
- dispose()
|
|
|
- mark(session, ['disposed'])
|
|
|
- expect(seen).toEqual(['test/marks'])
|
|
|
- })
|
|
|
-
|
|
|
- it('makes an earlier session listener read current and skips duplicate drive application', async () => {
|
|
|
- const ctx = new Context()
|
|
|
- await ctx.plugin(SessionStore)
|
|
|
- const seen: unknown[] = []
|
|
|
- ctx.on('session/event', (session) => {
|
|
|
- seen.push(ctx.sessionProjections.snapshot(session).values['test/marks'])
|
|
|
- })
|
|
|
- await ctx.plugin(SessionProjectionRegistry)
|
|
|
- ctx.sessionProjections.register(marksUnit())
|
|
|
- const session = ctx.sessions.create()
|
|
|
-
|
|
|
- mark(session, ['early listener'])
|
|
|
-
|
|
|
- expect(seen).toEqual([{ marks: ['early listener'] }])
|
|
|
- expect(ctx.sessionProjections.snapshot(session).values['test/marks'])
|
|
|
- .toEqual({ marks: ['early listener'] })
|
|
|
- })
|
|
|
-
|
|
|
- it('forward-applies events appended while the session is detached', async () => {
|
|
|
- const ctx = new Context()
|
|
|
- await ctx.plugin(SessionStore)
|
|
|
- await ctx.plugin(SessionProjectionRegistry)
|
|
|
- ctx.sessionProjections.register(countUnit())
|
|
|
- const session = ctx.sessions.prepare()
|
|
|
- const detach = ctx.sessions.enter(session)
|
|
|
- session.append('turn/start', { turn: 1 })
|
|
|
- expect(ctx.sessionProjections.snapshot(session).values['test/count']).toBe(1)
|
|
|
- detach()
|
|
|
- session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
|
|
|
- session.append('turn/start', { turn: 2 })
|
|
|
- ctx.sessions.enter(session)
|
|
|
-
|
|
|
- session.append('turn/end', { turn: 2, reason: { kind: 'completed' } })
|
|
|
-
|
|
|
- expect(ctx.sessionProjections.snapshot(session).values['test/count']).toBe(4)
|
|
|
- })
|
|
|
-
|
|
|
- it('brings a detached session cell current on its first read after reattachment', async () => {
|
|
|
- const ctx = new Context()
|
|
|
- await ctx.plugin(SessionStore)
|
|
|
- await ctx.plugin(SessionProjectionRegistry)
|
|
|
- ctx.sessionProjections.register(countUnit())
|
|
|
- const session = ctx.sessions.prepare()
|
|
|
- const detach = ctx.sessions.enter(session)
|
|
|
- session.append('turn/start', { turn: 1 })
|
|
|
- expect(ctx.sessionProjections.snapshot(session).values['test/count']).toBe(1)
|
|
|
- detach()
|
|
|
- session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
|
|
|
- session.append('turn/start', { turn: 2 })
|
|
|
- ctx.sessions.enter(session)
|
|
|
-
|
|
|
- expect(ctx.sessionProjections.snapshot(session).values['test/count']).toBe(3)
|
|
|
- })
|
|
|
-
|
|
|
it('drives independently per session (cells are per-session watermarks)', async () => {
|
|
|
const { ctx, session } = await harness()
|
|
|
const other = ctx.sessions.create()
|
|
|
@@ -213,9 +129,8 @@ describe('SessionProjectionRegistry drive', () => {
|
|
|
})
|
|
|
session.append('turn/start', { turn: 1 })
|
|
|
expect(changedKeys).toEqual([])
|
|
|
- const snapshot = ctx.sessionProjections.snapshot(session)
|
|
|
- expect(snapshot.values['test/count']).toBe(1)
|
|
|
- expect(snapshot.values['test/marks']).toEqual({ marks: [] })
|
|
|
+ expect(ctx.sessionProjections.stateOf(session, 'test/count')).toBe(1)
|
|
|
+ expect(ctx.sessionProjections.snapshot(session).values).toEqual({ 'test/marks': { marks: [] } })
|
|
|
})
|
|
|
|
|
|
it('shares one unit between registrants of the same key', async () => {
|
|
|
@@ -257,6 +172,14 @@ describe('SessionProjectionRegistry drive', () => {
|
|
|
.toThrow(/already registered at stateVersion 1; refusing to share it with stateVersion 9/)
|
|
|
})
|
|
|
|
|
|
+ it('refuses to share a key across a persistence-policy change', async () => {
|
|
|
+ const { ctx } = await harness()
|
|
|
+ ctx.sessionProjections.register(countUnit())
|
|
|
+
|
|
|
+ expect(() => ctx.sessionProjections.register({ ...countUnit(), persist: false }))
|
|
|
+ .toThrow(/already registered with persist true; refusing to share it with persist false/)
|
|
|
+ })
|
|
|
+
|
|
|
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/)
|
|
|
@@ -291,14 +214,15 @@ describe('SessionProjectionRegistry drive', () => {
|
|
|
expect(ctx.sessionProjections.snapshot(session).values).toEqual({})
|
|
|
})
|
|
|
|
|
|
- it('snapshot serves the wire view for a client key and the raw state for a host-internal key', async () => {
|
|
|
+ it('snapshot serves client views and excludes host-only state', async () => {
|
|
|
const { ctx, session } = await harness()
|
|
|
ctx.sessionProjections.register(marksUnit())
|
|
|
ctx.sessionProjections.register(countUnit())
|
|
|
mark(session, ['a', 'b'])
|
|
|
const values = ctx.sessionProjections.snapshot(session).values
|
|
|
expect(values['test/marks']).toEqual({ marks: ['a', 'b'] })
|
|
|
- expect(values['test/count']).toEqual(1)
|
|
|
+ expect('test/count' in values).toBe(false)
|
|
|
+ expect(ctx.sessionProjections.stateOf(session, 'test/count')).toBe(1)
|
|
|
expect('test/unregistered' in values).toBe(false)
|
|
|
})
|
|
|
|
|
|
@@ -379,7 +303,7 @@ describe('SessionProjectionRegistry drive', () => {
|
|
|
}, full, 0)
|
|
|
expect(snapshot.asOfSeq).toBe(4)
|
|
|
expect(snapshot.values['test/marks']).toEqual({ marks: ['new'] })
|
|
|
- expect(snapshot.values['test/count']).toBe(5) // refolded from init over all 5 events
|
|
|
+ expect('test/count' in snapshot.values).toBe(false)
|
|
|
// The refreshed rows sit at the served cut, ready for a durable write-back.
|
|
|
expect(checkpoint['test/marks']).toEqual({ ver: 1, seq: 4, val: { marks: ['new'] } })
|
|
|
expect(checkpoint['test/count']).toEqual({ ver: 1, seq: 4, val: 5 })
|
|
|
@@ -397,20 +321,22 @@ describe('SessionProjectionRegistry drive', () => {
|
|
|
{ type: 'turn/start', seq: 3, time: 3, data: { turn: 2 } },
|
|
|
{ type: 'turn/end', seq: 4, time: 4, data: { turn: 2, reason: { kind: 'completed' } } },
|
|
|
]
|
|
|
- const { snapshot } = ctx.sessionProjections.restore(rows, tail, 3)
|
|
|
+ const { snapshot, checkpoint } = ctx.sessionProjections.restore(rows, tail, 3)
|
|
|
expect(snapshot.asOfSeq).toBe(4)
|
|
|
// marks already covers the tail (watermark 4): nothing re-applied.
|
|
|
expect(snapshot.values['test/marks']).toEqual({ marks: ['done'] })
|
|
|
- // count folds exactly seqs 3 and 4 on top of its checkpoint.
|
|
|
- expect(snapshot.values['test/count']).toBe(5)
|
|
|
+ // count folds exactly seqs 3 and 4 on top of its checkpoint, but remains host-only.
|
|
|
+ expect(checkpoint['test/count']).toEqual({ ver: 1, seq: 4, val: 5 })
|
|
|
+ expect('test/count' in snapshot.values).toBe(false)
|
|
|
|
|
|
// Empty tail (checkpoint is current): the cut sits at baseSeq - 1.
|
|
|
- const { snapshot: current } = ctx.sessionProjections.restore({
|
|
|
+ const { snapshot: current, checkpoint: currentCheckpoint } = ctx.sessionProjections.restore({
|
|
|
'test/marks': { ver: 1, seq: 4, val: { marks: ['done'] } },
|
|
|
'test/count': { ver: 1, seq: 4, val: 5 },
|
|
|
}, [], 5)
|
|
|
expect(current.asOfSeq).toBe(4)
|
|
|
- expect(current.values['test/count']).toBe(5)
|
|
|
+ expect('test/count' in current.values).toBe(false)
|
|
|
+ expect(currentCheckpoint['test/count']).toEqual({ ver: 1, seq: 4, val: 5 })
|
|
|
})
|
|
|
|
|
|
it('viewCheckpoint serves version-matching rows without any log and skips mismatched keys', async () => {
|
|
|
@@ -426,7 +352,7 @@ describe('SessionProjectionRegistry drive', () => {
|
|
|
expect(ctx.sessionProjections.viewCheckpoint({})).toEqual({})
|
|
|
})
|
|
|
|
|
|
- it('viewCheckpoint and restore can serve wire keys while retaining full host checkpoints', async () => {
|
|
|
+ it('viewCheckpoint and restore exclude host-only state while retaining its checkpoint', async () => {
|
|
|
const { ctx } = await harness()
|
|
|
ctx.sessionProjections.register(marksUnit())
|
|
|
ctx.sessionProjections.register(countUnit())
|
|
|
@@ -436,13 +362,9 @@ describe('SessionProjectionRegistry drive', () => {
|
|
|
}
|
|
|
expect(ctx.sessionProjections.viewCheckpoint(rows)).toEqual({
|
|
|
'test/marks': { marks: ['stored'] },
|
|
|
- 'test/count': 5,
|
|
|
- })
|
|
|
- expect(ctx.sessionProjections.viewCheckpoint(rows, { wireOnly: true })).toEqual({
|
|
|
- 'test/marks': { marks: ['stored'] },
|
|
|
})
|
|
|
|
|
|
- const restored = ctx.sessionProjections.restore(rows, [], 5, { wireOnly: true })
|
|
|
+ const restored = ctx.sessionProjections.restore(rows, [], 5)
|
|
|
expect(restored.snapshot.values).toEqual({
|
|
|
'test/marks': { marks: ['stored'] },
|
|
|
})
|
|
|
@@ -470,7 +392,9 @@ describe('SessionProjectionRegistry drive', () => {
|
|
|
expect(floor).toBe(9)
|
|
|
// …an intact log serves the anchor event and the checkpoint stands as-is.
|
|
|
const anchor: SessionEvent = { type: 'turn/end', seq: 9, time: 9, data: { turn: 2, reason: { kind: 'completed' } } }
|
|
|
- expect(ctx.sessionProjections.restore(rows, [anchor], 9).snapshot.values['test/count']).toBe(10)
|
|
|
+ const anchored = ctx.sessionProjections.restore(rows, [anchor], 9)
|
|
|
+ expect(anchored.snapshot.values).toEqual({})
|
|
|
+ expect(anchored.checkpoint['test/count']).toEqual({ ver: 1, seq: 9, val: 10 })
|
|
|
// …while a log crash-repaired down to fewer events returns an empty tail:
|
|
|
// the row overreaches the proven end and a tail read cannot fix this key.
|
|
|
expect(() => ctx.sessionProjections.restore(rows, [], 9)).toThrow(/re-read from seq 0/)
|
|
|
@@ -479,9 +403,10 @@ describe('SessionProjectionRegistry drive', () => {
|
|
|
{ type: 'turn/start', seq: 0, time: 0, data: { turn: 1 } },
|
|
|
{ type: 'turn/end', seq: 1, time: 1, data: { turn: 1, reason: { kind: 'completed' } } },
|
|
|
]
|
|
|
- const { snapshot } = ctx.sessionProjections.restore(rows, events, 0)
|
|
|
+ const { snapshot, checkpoint } = ctx.sessionProjections.restore(rows, events, 0)
|
|
|
expect(snapshot.asOfSeq).toBe(1)
|
|
|
- expect(snapshot.values['test/count']).toBe(2)
|
|
|
+ expect(snapshot.values).toEqual({})
|
|
|
+ expect(checkpoint['test/count']).toEqual({ ver: 1, seq: 1, val: 2 })
|
|
|
})
|
|
|
|
|
|
it('fails loud when a unit view violates its own schema (async unit output is unrepresentable)', async () => {
|