registry.spec.ts 32 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777
  1. /**
  2. * SessionProjectionRegistry unit drive: eager apply on committed events with
  3. * lazy cell build (registration after events, session after registration),
  4. * the Object.is no-change gates (same state or raw view reference ⇒ zero
  5. * change-feed work), snapshot consistency (asOfSeq = last event seq; values
  6. * from the watermark cache), duplicate-key rejection, stateVersion validation,
  7. * and effect-tied removal of registrations and change listeners (HMR safety).
  8. */
  9. import { describe, expect, it, vi } from 'vitest'
  10. import { Context } from '@deepseek-ai/cordis'
  11. import { z } from 'zod'
  12. import SessionStore, {
  13. SESSION_FORMAT_VERSION,
  14. Session,
  15. SessionId,
  16. SessionLogOffset,
  17. SessionSeq,
  18. } from '@deepseek-ai/dsh-session'
  19. import type { SessionEvent, SessionHeader } from '@deepseek-ai/dsh-session'
  20. import SessionProjectionRegistry from '@deepseek-ai/dsh-session-projection'
  21. import type { ProjectionDefinition } from '@deepseek-ai/dsh-session-projection'
  22. declare module '@deepseek-ai/dsh-session-projection/types' {
  23. interface SessionProjectionStateMap {
  24. 'test/marks': MarksState
  25. 'test/count': number
  26. 'test/stable-view': StableViewState
  27. 'test/cut': number
  28. }
  29. interface SessionProjectionMap {
  30. 'test/marks': { marks: string[] }
  31. 'test/stable-view': { marks: string[] }
  32. }
  33. }
  34. declare module '@deepseek-ai/dsh-session/types' {
  35. interface SessionEventMap {
  36. 'test/mark': { marks: string[] }
  37. }
  38. }
  39. interface MarksView {
  40. marks: string[]
  41. }
  42. type MarksState = MarksView | null
  43. interface StableViewState {
  44. revision: number
  45. value: MarksView
  46. }
  47. const marksViewSchema: z.ZodType<MarksView> = z.object({ marks: z.array(z.string()) })
  48. const RESTORE_HEADER: SessionHeader = {
  49. version: SESSION_FORMAT_VERSION,
  50. id: SessionId('projection-restore'),
  51. createdAt: 0,
  52. isSeeded: false,
  53. }
  54. /** Whole-value unit: latest test/mark event wins; unrelated events return the same reference. */
  55. const marksUnit = (): Omit<ProjectionDefinition<'test/marks', MarksState>, 'wire'>
  56. & { wire: NonNullable<ProjectionDefinition<'test/marks', MarksState>['wire']> } => ({
  57. key: 'test/marks',
  58. stateSchema: marksViewSchema.nullable(),
  59. init: () => null,
  60. apply: (state, event) => (event.type === 'test/mark' ? (event).data : state),
  61. wire: {
  62. viewSchema: marksViewSchema,
  63. view: state => state ?? { marks: [] },
  64. },
  65. stateVersion: 1,
  66. })
  67. /** Host-only counting unit over every event — state changes on each apply. */
  68. const countUnit = (): ProjectionDefinition<'test/count', number> => ({
  69. key: 'test/count',
  70. stateSchema: z.number().int().nonnegative(),
  71. init: () => 0,
  72. apply: state => state + 1,
  73. stateVersion: 1,
  74. })
  75. const stableViewUnit = (
  76. view: (state: StableViewState) => StableViewState['value'],
  77. ) => ({
  78. key: 'test/stable-view',
  79. stateSchema: z.object({
  80. revision: z.number().int().nonnegative(),
  81. value: marksViewSchema,
  82. }),
  83. init: () => ({ revision: 0, value: { marks: [] } }),
  84. apply: (state, event) => {
  85. if (event.type === 'turn/start') return { ...state, revision: state.revision + 1 }
  86. if (event.type === 'test/mark') return { revision: state.revision + 1, value: event.data }
  87. return state
  88. },
  89. wire: {
  90. viewSchema: marksViewSchema,
  91. view,
  92. },
  93. stateVersion: 1,
  94. }) satisfies ProjectionDefinition<'test/stable-view', StableViewState>
  95. /** Host-only unit whose initial state proves the exact inherited cut. */
  96. const cutUnit = (): ProjectionDefinition<'test/cut', number> => ({
  97. key: 'test/cut',
  98. stateSchema: z.number().int().nonnegative(),
  99. init: (_header, inheritedEventCount) => inheritedEventCount,
  100. apply: state => state,
  101. stateVersion: 1,
  102. })
  103. async function harness(): Promise<{ ctx: Context; session: Session }> {
  104. const ctx = new Context()
  105. await ctx.plugin(SessionStore)
  106. await ctx.plugin(SessionProjectionRegistry)
  107. return { ctx, session: ctx.sessions.create() }
  108. }
  109. const mark = (session: Session, marks: string[]): SessionEvent =>
  110. session.append('test/mark', { marks })
  111. const STATE_SEQUENCES = [
  112. [0, 0, 0, 0],
  113. [0, 0, 0, 1],
  114. [0, 0, 1, 0],
  115. [0, 0, 1, 1],
  116. [0, 0, 1, 2],
  117. [0, 1, 0, 0],
  118. [0, 1, 0, 1],
  119. [0, 1, 0, 2],
  120. [0, 1, 1, 0],
  121. [0, 1, 1, 1],
  122. [0, 1, 1, 2],
  123. [0, 1, 2, 0],
  124. [0, 1, 2, 1],
  125. [0, 1, 2, 2],
  126. [0, 1, 2, 3],
  127. ] as const
  128. function identitySequences(length: number): number[][] {
  129. const sequences: number[][] = []
  130. const visit = (sequence: number[], highest: number): void => {
  131. if (sequence.length === length) {
  132. sequences.push(sequence)
  133. return
  134. }
  135. for (let value = 0; value <= highest + 1; value++) {
  136. visit([...sequence, value], Math.max(highest, value))
  137. }
  138. }
  139. visit([0], 0)
  140. return sequences
  141. }
  142. function sameIdentities(left: readonly unknown[], right: readonly unknown[]): boolean {
  143. return left.length === right.length && left.every((value, index) => Object.is(value, right[index]))
  144. }
  145. function sequenceName(sequence: readonly number[], prefix: string): string {
  146. return sequence.map(value => `${prefix}${String(value + 1)}`).join(',')
  147. }
  148. describe('SessionProjectionRegistry drive', () => {
  149. it('supplies the exact inherited cut to live, restored, and hydrated projection initialization', async () => {
  150. const ctx = new Context()
  151. await ctx.plugin(SessionStore)
  152. await ctx.plugin(SessionProjectionRegistry)
  153. ctx.sessionProjections.register(cutUnit())
  154. const inherited: SessionEvent[] = [
  155. { type: 'turn/start', seq: SessionSeq(0), time: 1, data: { turn: 1 } },
  156. { type: 'turn/end', seq: SessionSeq(1), time: 2, data: { turn: 1, reason: { kind: 'completed' } } },
  157. ]
  158. const session = ctx.sessions.create(SessionId('projection-cut'), {
  159. seed: inherited,
  160. inheritedEventCount: SessionLogOffset(inherited.length),
  161. meta: { isSeeded: true },
  162. })
  163. expect(ctx.sessionProjections.stateOf(session, 'test/cut')).toBe(inherited.length)
  164. const restored = ctx.sessionProjections.restore(
  165. {},
  166. inherited,
  167. SessionLogOffset(0),
  168. session.header,
  169. session.inheritedEventCount,
  170. )
  171. expect(restored.checkpoint['test/cut']?.val).toBe(inherited.length)
  172. const prepared = Session.create(
  173. SessionId('projection-cut-prepared'),
  174. inherited,
  175. { ...session.header, id: SessionId('projection-cut-prepared') },
  176. session.inheritedEventCount,
  177. )
  178. expect(ctx.sessionProjections.hydrate(
  179. prepared,
  180. {},
  181. inherited,
  182. SessionLogOffset(0),
  183. ).asOfSeq).toBe(1)
  184. expect(ctx.sessionProjections.stateOf(prepared, 'test/cut')).toBe(inherited.length)
  185. })
  186. it('drives a registered unit over committed events and snapshots the current value', async () => {
  187. const { ctx, session } = await harness()
  188. ctx.sessionProjections.register(marksUnit())
  189. mark(session, ['a'])
  190. mark(session, ['a', 'b'])
  191. const snapshot = ctx.sessionProjections.snapshot(session)
  192. expect(snapshot.values['test/marks']).toEqual({ marks: ['a', 'b'] })
  193. expect(snapshot.asOfSeq).toBe(session.seq - 1)
  194. })
  195. it('builds the cell lazily from the full log for a unit registered after events flowed', async () => {
  196. const { ctx, session } = await harness()
  197. mark(session, ['pre-registration'])
  198. ctx.sessionProjections.register(marksUnit())
  199. expect(ctx.sessionProjections.snapshot(session).values['test/marks']).toEqual({ marks: ['pre-registration'] })
  200. // The lazily-built cell then continues on the live drive path.
  201. mark(session, ['after'])
  202. expect(ctx.sessionProjections.snapshot(session).values['test/marks']).toEqual({ marks: ['after'] })
  203. })
  204. it('serves init-derived state and asOfSeq -1 for an empty log', async () => {
  205. const { ctx, session } = await harness()
  206. ctx.sessionProjections.register(marksUnit())
  207. const snapshot = ctx.sessionProjections.snapshot(session)
  208. expect(snapshot.asOfSeq).toBe(-1)
  209. expect(snapshot.values['test/marks']).toEqual({ marks: [] })
  210. })
  211. it('notifies onChanged with the validated view and the causing seq, and skips same-reference applies', async () => {
  212. const { ctx, session } = await harness()
  213. ctx.sessionProjections.register(marksUnit())
  214. const seen: { key: string; value: unknown; seq: SessionSeq; sessionId: string }[] = []
  215. ctx.sessionProjections.onChanged((changedSession, key, value, seq) => {
  216. seen.push({ key, value, seq, sessionId: String(changedSession.id) })
  217. })
  218. const event = mark(session, ['a'])
  219. // Non-matching event: apply returns the same reference — no notification.
  220. session.append('turn/start', { turn: 1 })
  221. expect(seen).toEqual([{ key: 'test/marks', value: { marks: ['a'] }, seq: event.seq, sessionId: String(session.id) }])
  222. })
  223. it('does not compute a view while no change listener exists', async () => {
  224. const { ctx, session } = await harness()
  225. const view = vi.fn((state: StableViewState) => state.value)
  226. ctx.sessionProjections.register(stableViewUnit(view))
  227. session.append('turn/start', { turn: 1 })
  228. session.append('turn/start', { turn: 2 })
  229. expect(ctx.sessionProjections.stateOf(session, 'test/stable-view')?.revision).toBe(2)
  230. expect(view).not.toHaveBeenCalled()
  231. })
  232. it('publishes the first observed view and suppresses later same-reference views', async () => {
  233. const { ctx, session } = await harness()
  234. const view = vi.fn((state: StableViewState) => state.value)
  235. ctx.sessionProjections.register(stableViewUnit(view))
  236. const seen: unknown[] = []
  237. ctx.sessionProjections.onChanged((_session, key, value) => {
  238. if (key === 'test/stable-view') seen.push(value)
  239. })
  240. session.append('turn/start', { turn: 1 })
  241. session.append('turn/start', { turn: 2 })
  242. expect(seen).toEqual([{ marks: [] }])
  243. expect(view).toHaveBeenCalledTimes(2)
  244. mark(session, ['changed'])
  245. expect(seen).toEqual([{ marks: [] }, { marks: ['changed'] }])
  246. expect(view).toHaveBeenCalledTimes(3)
  247. })
  248. it('publishes the first view after an unobserved state change', async () => {
  249. const { ctx, session } = await harness()
  250. const view = vi.fn((state: StableViewState) => state.value)
  251. ctx.sessionProjections.register(stableViewUnit(view))
  252. const first: unknown[] = []
  253. const stop = ctx.sessionProjections.onChanged((_session, key, value) => {
  254. if (key === 'test/stable-view') first.push(value)
  255. })
  256. session.append('turn/start', { turn: 1 })
  257. stop()
  258. session.append('turn/start', { turn: 2 })
  259. expect(view).toHaveBeenCalledTimes(1)
  260. const resumed: unknown[] = []
  261. ctx.sessionProjections.onChanged((_session, key, value) => {
  262. if (key === 'test/stable-view') resumed.push(value)
  263. })
  264. session.append('turn/start', { turn: 3 })
  265. expect(first).toEqual([{ marks: [] }])
  266. expect(resumed).toEqual([{ marks: [] }])
  267. expect(view).toHaveBeenCalledTimes(2)
  268. })
  269. it('matches every four-state identity sequence across listener gaps and raw-view identities', async () => {
  270. const { ctx } = await harness()
  271. const initialState: MarksState = { marks: ['initial'] }
  272. const stateByEvent = new Map<string, MarksState>()
  273. const viewByState = new Map<MarksState, MarksView>()
  274. const computedViews: MarksView[] = []
  275. ctx.sessionProjections.register({
  276. key: 'test/marks',
  277. stateSchema: marksViewSchema.nullable(),
  278. init: () => initialState,
  279. apply: (state, event) => {
  280. if (event.type !== 'test/mark') return state
  281. const token = event.data.marks[0]
  282. if (token === undefined || !stateByEvent.has(token)) return state
  283. return stateByEvent.get(token) as MarksState
  284. },
  285. wire: {
  286. viewSchema: marksViewSchema,
  287. view: (state) => {
  288. const value = viewByState.get(state)
  289. if (value === undefined) throw new Error('test state lacks a raw view')
  290. computedViews.push(value)
  291. return value
  292. },
  293. },
  294. stateVersion: 1,
  295. })
  296. const failures = new Map<string, unknown>()
  297. let mismatchCount = 0
  298. let checked = 0
  299. for (const stateSequence of STATE_SEQUENCES) {
  300. const stateCount = Math.max(...stateSequence) + 1
  301. for (const viewSequence of identitySequences(stateCount)) {
  302. for (const baselineKnown of [false, true]) {
  303. for (let listenerMask = 0; listenerMask < 8; listenerMask++) {
  304. const scenario = String(checked++)
  305. const states = Array.from(
  306. { length: stateCount },
  307. (_, index): MarksState => ({ marks: [`state-${scenario}-${String(index)}`] }),
  308. )
  309. const views = Array.from(
  310. { length: Math.max(...viewSequence) + 1 },
  311. (): MarksView => ({ marks: [] }),
  312. )
  313. for (let index = 0; index < stateCount; index++) {
  314. viewByState.set(states[index] as MarksState, views[viewSequence[index] as number] as MarksView)
  315. }
  316. for (let index = 0; index < stateSequence.length; index++) {
  317. stateByEvent.set(`${scenario}:${String(index)}`, states[stateSequence[index] as number] as MarksState)
  318. }
  319. const session = ctx.sessions.create()
  320. const notifications: number[] = []
  321. let stop: (() => void) | undefined
  322. const setListening = (listening: boolean): void => {
  323. if (listening && stop === undefined) {
  324. stop = ctx.sessionProjections.onChanged((changedSession, key, _value, seq) => {
  325. if (changedSession === session && key === 'test/marks') notifications.push(seq)
  326. })
  327. } else if (!listening && stop !== undefined) {
  328. stop()
  329. stop = undefined
  330. }
  331. }
  332. setListening(baselineKnown)
  333. mark(session, [`${scenario}:0`])
  334. computedViews.length = 0
  335. notifications.length = 0
  336. const expectedViews: MarksView[] = []
  337. const expectedNotifications: number[] = []
  338. let comparable = baselineKnown
  339. ? views[viewSequence[stateSequence[0] as number] as number] as MarksView
  340. : undefined
  341. for (let index = 1; index < stateSequence.length; index++) {
  342. const listening = (listenerMask & (1 << (index - 1))) !== 0
  343. setListening(listening)
  344. const changed = stateSequence[index] !== stateSequence[index - 1]
  345. if (changed) {
  346. if (listening) {
  347. const current = views[viewSequence[stateSequence[index] as number] as number] as MarksView
  348. expectedViews.push(current)
  349. if (comparable === undefined || !Object.is(comparable, current)) {
  350. expectedNotifications.push(index)
  351. }
  352. comparable = current
  353. } else {
  354. comparable = undefined
  355. }
  356. }
  357. mark(session, [`${scenario}:${String(index)}`])
  358. }
  359. setListening(false)
  360. if (!sameIdentities(computedViews, expectedViews)
  361. || notifications.length !== expectedNotifications.length
  362. || notifications.some((seq, index) => seq !== expectedNotifications[index])) {
  363. mismatchCount += 1
  364. const stateName = sequenceName(stateSequence, 'v')
  365. if (!failures.has(stateName) || (baselineKnown && listenerMask === 7)) {
  366. failures.set(stateName, {
  367. state: stateName,
  368. view: stateSequence.map(value => `r${String((viewSequence[value] as number) + 1)}`).join(','),
  369. baseline: baselineKnown ? 'known' : 'unknown',
  370. listeners: [0, 1, 2]
  371. .map(index => (listenerMask & (1 << index)) === 0 ? 'off' : 'on')
  372. .join(','),
  373. expectedViewCalls: expectedViews.length,
  374. actualViewCalls: computedViews.length,
  375. expectedNotifications,
  376. actualNotifications: [...notifications],
  377. })
  378. }
  379. }
  380. computedViews.length = 0
  381. }
  382. }
  383. }
  384. }
  385. expect({ checked, mismatchCount, failures: [...failures.values()] }).toEqual({
  386. checked: 960,
  387. mismatchCount: 0,
  388. failures: [],
  389. })
  390. })
  391. it('drives independently per session (cells are per-session watermarks)', async () => {
  392. const { ctx, session } = await harness()
  393. const other = ctx.sessions.create()
  394. ctx.sessionProjections.register(marksUnit())
  395. mark(session, ['one'])
  396. mark(other, ['two'])
  397. expect(ctx.sessionProjections.snapshot(session).values['test/marks']).toEqual({ marks: ['one'] })
  398. expect(ctx.sessionProjections.snapshot(other).values['test/marks']).toEqual({ marks: ['two'] })
  399. })
  400. it('updates host-only units without publishing them to wire listeners', async () => {
  401. const { ctx, session } = await harness()
  402. ctx.sessionProjections.register(marksUnit())
  403. ctx.sessionProjections.register(countUnit())
  404. const changedKeys: string[] = []
  405. ctx.sessionProjections.onChanged((_session, key) => {
  406. changedKeys.push(key)
  407. })
  408. session.append('turn/start', { turn: 1 })
  409. expect(changedKeys).toEqual([])
  410. expect(ctx.sessionProjections.stateOf(session, 'test/count')).toBe(1)
  411. expect(ctx.sessionProjections.snapshot(session).values).toEqual({ 'test/marks': { marks: [] } })
  412. })
  413. it('shares one unit between registrants of the same key', async () => {
  414. const { ctx, session } = await harness()
  415. ctx.sessionProjections.register(marksUnit())
  416. // One definition already serves every session (cells are keyed by
  417. // Session), and registrants are per-session now: an agent preset mounts
  418. // the same tool package once per agent.
  419. expect(() => ctx.sessionProjections.register(marksUnit())).not.toThrow()
  420. mark(session, ['kept'])
  421. expect(ctx.sessionProjections.snapshot(session).values['test/marks']).toEqual({ marks: ['kept'] })
  422. })
  423. it('keeps the unit until the last registrant releases it', async () => {
  424. const { ctx, session } = await harness()
  425. const first = ctx.sessionProjections.register(marksUnit())
  426. const second = ctx.sessionProjections.register(marksUnit())
  427. mark(session, ['kept'])
  428. first()
  429. // The regression this counts against: without last-release semantics, one
  430. // session ending strips the projection from every other live session,
  431. // because the first registrant owns the only disposer.
  432. expect(ctx.sessionProjections.snapshot(session).values['test/marks']).toEqual({ marks: ['kept'] })
  433. second()
  434. expect(ctx.sessionProjections.snapshot(session).values).toEqual({})
  435. })
  436. it('refuses to share a key across a stateVersion change', async () => {
  437. const { ctx } = await harness()
  438. ctx.sessionProjections.register(marksUnit())
  439. // The one incompatibility a runtime comparison can name: the versioned
  440. // contract says the cached state shape differs, so the two cannot share
  441. // cells. Everything else about a definition is functions.
  442. expect(() => ctx.sessionProjections.register({ ...marksUnit(), stateVersion: 9 }))
  443. .toThrow(/already registered at stateVersion 1; refusing to share it with stateVersion 9/)
  444. })
  445. it('rejects a non-integer or negative stateVersion at register time', async () => {
  446. const { ctx } = await harness()
  447. expect(() => ctx.sessionProjections.register({ ...marksUnit(), stateVersion: -1 })).toThrow(/stateVersion/)
  448. expect(() => ctx.sessionProjections.register({ ...marksUnit(), stateVersion: 1.5 })).toThrow(/stateVersion/)
  449. })
  450. it('register() disposer removes the key (with its cells) and frees it for re-registration', async () => {
  451. const { ctx, session } = await harness()
  452. const dispose = ctx.sessionProjections.register(marksUnit())
  453. mark(session, ['cached'])
  454. dispose()
  455. expect(ctx.sessionProjections.snapshot(session).values).toEqual({})
  456. ctx.sessionProjections.register(marksUnit())
  457. // Fresh registration rebuilds from the log, not from a stale cell.
  458. expect(ctx.sessionProjections.snapshot(session).values['test/marks']).toEqual({ marks: ['cached'] })
  459. })
  460. it('removes registrations and change listeners when their owning fiber unloads (HMR safety)', async () => {
  461. const { ctx, session } = await harness()
  462. const notifications: string[] = []
  463. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  464. inner.sessionProjections.register(marksUnit())
  465. inner.sessionProjections.onChanged((_session, key) => {
  466. notifications.push(key)
  467. })
  468. }, { inject: ['sessionProjections'] }))
  469. mark(session, ['live'])
  470. expect(notifications).toEqual(['test/marks'])
  471. await fiber.dispose()
  472. mark(session, ['after-dispose'])
  473. expect(notifications).toEqual(['test/marks'])
  474. expect(ctx.sessionProjections.snapshot(session).values).toEqual({})
  475. })
  476. it('snapshot serves client views and excludes host-only state', async () => {
  477. const { ctx, session } = await harness()
  478. ctx.sessionProjections.register(marksUnit())
  479. ctx.sessionProjections.register(countUnit())
  480. mark(session, ['a', 'b'])
  481. const values = ctx.sessionProjections.snapshot(session).values
  482. expect(values['test/marks']).toEqual({ marks: ['a', 'b'] })
  483. expect('test/count' in values).toBe(false)
  484. expect(ctx.sessionProjections.stateOf(session, 'test/count')).toBe(1)
  485. expect('test/unregistered' in values).toBe(false)
  486. })
  487. it('checkpoints every persisted unit with its stateVersion and per-cell watermark', async () => {
  488. const { ctx, session } = await harness()
  489. ctx.sessionProjections.register(marksUnit())
  490. ctx.sessionProjections.register({ ...countUnit(), stateVersion: 7 })
  491. const markEvent = mark(session, ['a'])
  492. const rows = ctx.sessionProjections.checkpoint(session)
  493. expect(rows['test/marks']).toEqual({ ver: 1, seq: markEvent.seq, val: { marks: ['a'] } })
  494. expect(rows['test/count']).toEqual({ ver: 7, seq: markEvent.seq, val: 1 })
  495. // Empty log: init-derived state at watermark -1.
  496. const fresh = ctx.sessions.create()
  497. expect(ctx.sessionProjections.checkpoint(fresh)['test/marks']).toEqual({ ver: 1, seq: -1, val: null })
  498. })
  499. it('checkpoint states are detached clones — mutating them cannot corrupt the watermark cache', async () => {
  500. const { ctx, session } = await harness()
  501. ctx.sessionProjections.register(marksUnit())
  502. mark(session, ['a'])
  503. const rows = ctx.sessionProjections.checkpoint(session)
  504. // Hostile (or merely careless) consumer mutates the handed-out state.
  505. ;(rows['test/marks']?.val as { marks: string[] }).marks.push('INJECTED')
  506. // The registry's authoritative cell is untouched: snapshot and a fresh
  507. // checkpoint both still serve the committed value.
  508. expect(ctx.sessionProjections.snapshot(session).values['test/marks']).toEqual({ marks: ['a'] })
  509. expect(ctx.sessionProjections.checkpoint(session)['test/marks']?.val).toEqual({ marks: ['a'] })
  510. })
  511. it('restoreFloor anchors one below the lowest usable watermark and at 0 for missing or mismatched rows', async () => {
  512. const { ctx } = await harness()
  513. expect(ctx.sessionProjections.restoreFloor({})).toBeUndefined() // no unit registered
  514. ctx.sessionProjections.register(marksUnit())
  515. ctx.sessionProjections.register(countUnit())
  516. expect(ctx.sessionProjections.restoreFloor({})).toBe(0)
  517. // Lowest usable watermark is count's 5 → the anchored tail starts AT 5
  518. // (one below the first needed seq 6), so the read proves seq 5 still exists.
  519. expect(ctx.sessionProjections.restoreFloor({
  520. 'test/marks': { ver: 1, seq: SessionSeq(10), val: { marks: [] } },
  521. 'test/count': { ver: 1, seq: SessionSeq(5), val: 6 },
  522. })).toBe(5)
  523. // A version-mismatched row forces that key back to a full refold.
  524. expect(ctx.sessionProjections.restoreFloor({
  525. 'test/marks': { ver: 2, seq: SessionSeq(10), val: { marks: [] } },
  526. 'test/count': { ver: 1, seq: SessionSeq(5), val: 6 },
  527. })).toBe(0)
  528. // A fresh (-1) row still needs the whole tail from 0.
  529. expect(ctx.sessionProjections.restoreFloor({
  530. 'test/marks': { ver: 1, seq: -1, val: null },
  531. 'test/count': { ver: 1, seq: -1, val: 0 },
  532. })).toBe(0)
  533. })
  534. it('restore folds the tail past each usable row and refolds from init on version mismatch', async () => {
  535. const { ctx } = await harness()
  536. ctx.sessionProjections.register(marksUnit())
  537. ctx.sessionProjections.register(countUnit())
  538. const tail: SessionEvent[] = [
  539. { type: 'test/mark', seq: SessionSeq(3), time: 3, data: { marks: ['new'] } },
  540. { type: 'turn/end', seq: SessionSeq(4), time: 4, data: { turn: 1, reason: { kind: 'completed' } } },
  541. ]
  542. // marks row usable (watermark 2, tail starts at 3); count row mismatched — but
  543. // a mismatch with baseSeq > 0 cannot silently refold: it throws for a re-read.
  544. expect(() => ctx.sessionProjections.restore({
  545. 'test/marks': { ver: 1, seq: SessionSeq(2), val: { marks: ['old'] } },
  546. 'test/count': { ver: 99, seq: SessionSeq(2), val: 3 },
  547. }, tail, SessionLogOffset(3), RESTORE_HEADER, SessionLogOffset(0)))
  548. .toThrow(/re-read from seq 0/)
  549. // The full-log re-read (baseSeq 0) refolds the mismatched key from init.
  550. const full: SessionEvent[] = [
  551. { type: 'turn/start', seq: SessionSeq(0), time: 0, data: { turn: 1 } },
  552. { type: 'test/mark', seq: SessionSeq(1), time: 1, data: { marks: ['old'] } },
  553. { type: 'test/mark', seq: SessionSeq(2), time: 2, data: { marks: ['old', '2'] } },
  554. ...tail,
  555. ]
  556. const { snapshot, checkpoint } = ctx.sessionProjections.restore({
  557. 'test/marks': { ver: 1, seq: SessionSeq(2), val: { marks: ['old', '2'] } },
  558. 'test/count': { ver: 99, seq: SessionSeq(2), val: 3 },
  559. }, full, SessionLogOffset(0), RESTORE_HEADER, SessionLogOffset(0))
  560. expect(snapshot.asOfSeq).toBe(4)
  561. expect(snapshot.values['test/marks']).toEqual({ marks: ['new'] })
  562. expect('test/count' in snapshot.values).toBe(false)
  563. // The refreshed rows sit at the served cut, ready for a durable write-back.
  564. expect(checkpoint['test/marks']).toEqual({ ver: 1, seq: 4, val: { marks: ['new'] } })
  565. expect(checkpoint['test/count']).toEqual({ ver: 1, seq: 4, val: 5 })
  566. })
  567. it('restore over a suffix folds only past each row watermark and serves an exact empty-tail cut', async () => {
  568. const { ctx } = await harness()
  569. ctx.sessionProjections.register(marksUnit())
  570. ctx.sessionProjections.register(countUnit())
  571. const rows = {
  572. 'test/marks': { ver: 1, seq: SessionSeq(4), val: { marks: ['done'] } },
  573. 'test/count': { ver: 1, seq: SessionSeq(2), val: 3 },
  574. }
  575. const tail: SessionEvent[] = [
  576. { type: 'turn/start', seq: SessionSeq(3), time: 3, data: { turn: 2 } },
  577. { type: 'turn/end', seq: SessionSeq(4), time: 4, data: { turn: 2, reason: { kind: 'completed' } } },
  578. ]
  579. const { snapshot, checkpoint } = ctx.sessionProjections.restore(
  580. rows,
  581. tail,
  582. SessionLogOffset(3),
  583. RESTORE_HEADER,
  584. SessionLogOffset(0),
  585. )
  586. expect(snapshot.asOfSeq).toBe(4)
  587. // marks already covers the tail (watermark 4): nothing re-applied.
  588. expect(snapshot.values['test/marks']).toEqual({ marks: ['done'] })
  589. // count folds exactly seqs 3 and 4 on top of its checkpoint, but remains host-only.
  590. expect(checkpoint['test/count']).toEqual({ ver: 1, seq: 4, val: 5 })
  591. expect('test/count' in snapshot.values).toBe(false)
  592. // Empty tail (checkpoint is current): the cut sits at baseSeq - 1.
  593. const { snapshot: current, checkpoint: currentCheckpoint } = ctx.sessionProjections.restore({
  594. 'test/marks': { ver: 1, seq: SessionSeq(4), val: { marks: ['done'] } },
  595. 'test/count': { ver: 1, seq: SessionSeq(4), val: 5 },
  596. }, [], SessionLogOffset(5), RESTORE_HEADER, SessionLogOffset(0))
  597. expect(current.asOfSeq).toBe(4)
  598. expect('test/count' in current.values).toBe(false)
  599. expect(currentCheckpoint['test/count']).toEqual({ ver: 1, seq: 4, val: 5 })
  600. })
  601. it('viewCheckpoint serves version-matching rows without any log and skips mismatched keys', async () => {
  602. const { ctx } = await harness()
  603. ctx.sessionProjections.register(marksUnit())
  604. ctx.sessionProjections.register(countUnit())
  605. const values = ctx.sessionProjections.viewCheckpoint({
  606. 'test/marks': { ver: 1, seq: SessionSeq(4), val: { marks: ['stored'] } },
  607. 'test/count': { ver: 99, seq: SessionSeq(4), val: 5 }, // mismatched: absent
  608. })
  609. expect(values['test/marks']).toEqual({ marks: ['stored'] })
  610. expect('test/count' in values).toBe(false)
  611. expect(ctx.sessionProjections.viewCheckpoint({})).toEqual({})
  612. })
  613. it('viewCheckpoint and restore exclude host-only state while retaining its checkpoint', async () => {
  614. const { ctx } = await harness()
  615. ctx.sessionProjections.register(marksUnit())
  616. ctx.sessionProjections.register(countUnit())
  617. const rows = {
  618. 'test/marks': { ver: 1, seq: SessionSeq(4), val: { marks: ['stored'] } },
  619. 'test/count': { ver: 1, seq: SessionSeq(4), val: 5 },
  620. }
  621. expect(ctx.sessionProjections.viewCheckpoint(rows)).toEqual({
  622. 'test/marks': { marks: ['stored'] },
  623. })
  624. const restored = ctx.sessionProjections.restore(
  625. rows,
  626. [],
  627. SessionLogOffset(5),
  628. RESTORE_HEADER,
  629. SessionLogOffset(0),
  630. )
  631. expect(restored.snapshot.values).toEqual({
  632. 'test/marks': { marks: ['stored'] },
  633. })
  634. expect(restored.checkpoint['test/count']).toEqual(rows['test/count'])
  635. })
  636. it('rejects version-matching rows whose state no longer matches the registered schema', async () => {
  637. const { ctx } = await harness()
  638. ctx.sessionProjections.register(marksUnit())
  639. const drifted = {
  640. 'test/marks': { ver: 1, seq: SessionSeq(2), val: { marks: 'not-an-array' } },
  641. }
  642. expect(ctx.sessionProjections.viewCheckpoint(drifted)).toEqual({})
  643. expect(() => ctx.sessionProjections.restore(
  644. drifted,
  645. [],
  646. SessionLogOffset(3),
  647. RESTORE_HEADER,
  648. SessionLogOffset(0),
  649. )).toThrow()
  650. })
  651. it('restore rejects a row claiming events past the supplied log end (shrunk log ⇒ re-read)', async () => {
  652. const { ctx } = await harness()
  653. ctx.sessionProjections.register(countUnit())
  654. const rows = { 'test/count': { ver: 1, seq: SessionSeq(9), val: 10 } }
  655. // The anchored floor sits ON the watermark, so the tail read must return
  656. // at least seq 9 from an intact log…
  657. const floor = ctx.sessionProjections.restoreFloor(rows)
  658. expect(floor).toBe(9)
  659. // …an intact log serves the anchor event and the checkpoint stands as-is.
  660. const anchor: SessionEvent = { type: 'turn/end', seq: SessionSeq(9), time: 9, data: { turn: 2, reason: { kind: 'completed' } } }
  661. const anchored = ctx.sessionProjections.restore(
  662. rows,
  663. [anchor],
  664. SessionLogOffset(9),
  665. RESTORE_HEADER,
  666. SessionLogOffset(0),
  667. )
  668. expect(anchored.snapshot.values).toEqual({})
  669. expect(anchored.checkpoint['test/count']).toEqual({ ver: 1, seq: 9, val: 10 })
  670. // …while a log crash-repaired down to fewer events returns an empty tail:
  671. // the row overreaches the proven end and a tail read cannot fix this key.
  672. expect(() => ctx.sessionProjections.restore(
  673. rows,
  674. [],
  675. SessionLogOffset(9),
  676. RESTORE_HEADER,
  677. SessionLogOffset(0),
  678. )).toThrow(/re-read from seq 0/)
  679. // The full re-read discards the overreaching row and refolds from init.
  680. const events: SessionEvent[] = [
  681. { type: 'turn/start', seq: SessionSeq(0), time: 0, data: { turn: 1 } },
  682. { type: 'turn/end', seq: SessionSeq(1), time: 1, data: { turn: 1, reason: { kind: 'completed' } } },
  683. ]
  684. const { snapshot, checkpoint } = ctx.sessionProjections.restore(
  685. rows,
  686. events,
  687. SessionLogOffset(0),
  688. RESTORE_HEADER,
  689. SessionLogOffset(0),
  690. )
  691. expect(snapshot.asOfSeq).toBe(1)
  692. expect(snapshot.values).toEqual({})
  693. expect(checkpoint['test/count']).toEqual({ ver: 1, seq: 1, val: 2 })
  694. })
  695. it('fails loud when a unit view violates its own schema (async unit output is unrepresentable)', async () => {
  696. const { ctx, session } = await harness()
  697. ctx.sessionProjections.register({
  698. key: 'test/marks',
  699. stateSchema: z.object({ marks: z.array(z.string()) }).nullable(),
  700. init: () => null as MarksState,
  701. apply: state => state,
  702. wire: {
  703. viewSchema: z.object({ marks: z.array(z.string()) }),
  704. // A Promise (what an accidentally-async view would return) is not the
  705. // declared shape: the boundary parse rejects it before it leaves.
  706. view: () => Promise.resolve({ marks: [] }) as never,
  707. },
  708. stateVersion: 1,
  709. })
  710. expect(() => ctx.sessionProjections.snapshot(session)).toThrow()
  711. })
  712. })