cache.spec.ts 25 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573
  1. /**
  2. * SessionProjectionCache behavior: mandatory-point writes (turn/end, detach),
  3. * count/interval throttling between them, fail-soft durability (a failed
  4. * write logs and stays stale, never throws into the event path), and the
  5. * synchronous cached listing read. The durable medium is the
  6. * `session_projcache` storage domain in per-record layout: one
  7. * version-stamped document per session under the json backend root at
  8. * `<root>/session_projcache/sessions/<id>.json`. Reads never touch the
  9. * medium — they come from the domain's in-memory tables, which writes mutate
  10. * only after durability.
  11. */
  12. import { afterEach, describe, expect, it, vi } from 'vitest'
  13. import { mkdir, mkdtemp, readFile, rm, writeFile } from 'node:fs/promises'
  14. import { tmpdir } from 'node:os'
  15. import { dirname, join } from 'node:path'
  16. import { Context } from '@deepseek-ai/cordis'
  17. import { z } from 'zod'
  18. import SessionStore, {
  19. Session,
  20. SessionId,
  21. SessionLogOffset,
  22. SessionSeq,
  23. } from '@deepseek-ai/dsh-session'
  24. import type { SessionEvent } from '@deepseek-ai/dsh-session'
  25. import SessionProjectionRegistry from '@deepseek-ai/dsh-session-projection'
  26. import type { ProjectionDefinition } from '@deepseek-ai/dsh-session-projection'
  27. import Storage from '@deepseek-ai/dsh-storage'
  28. import {
  29. apply as storageJsonApply, Config as storageJsonConfig, inject as storageJsonInject, name as storageJsonName,
  30. } from '@deepseek-ai/dsh-storage-json'
  31. import {
  32. apply as storageDomainApply, Config as storageDomainConfig, inject as storageDomainInject, name as storageDomainName,
  33. } from '@deepseek-ai/dsh-storage-domain'
  34. import SessionProjectionCache from '../src/index.ts'
  35. import { checkpointRecord, projectionCacheDomainSpec } from '../src/spec.ts'
  36. import type { CheckpointRecord } from '../src/spec.ts'
  37. declare module '@deepseek-ai/dsh-session-projection/types' {
  38. interface SessionProjectionStateMap {
  39. 'cache-test/marks': MarksState
  40. 'cache-test/marks2': Map<string, string>
  41. 'cache-test/count': number
  42. 'cache-test/secret': string
  43. }
  44. interface SessionProjectionMap {
  45. 'cache-test/marks': { marks: string[] }
  46. }
  47. }
  48. declare module '@deepseek-ai/dsh-session/types' {
  49. interface SessionEventMap {
  50. 'cache-test/mark': { marks: string[] }
  51. }
  52. interface OutOfBandSessionEventMap {
  53. 'cache-test/mark': true
  54. }
  55. }
  56. type MarksState = { marks: string[] } | null
  57. const marksUnit = (stateVersion = 1) => ({
  58. key: 'cache-test/marks',
  59. stateSchema: z.object({ marks: z.array(z.string()) }).nullable(),
  60. init: () => null,
  61. apply: (state, event) => (event.type === 'cache-test/mark' ? (event).data : state),
  62. wire: {
  63. viewSchema: z.object({ marks: z.array(z.string()) }),
  64. view: state => state ?? { marks: [] },
  65. },
  66. stateVersion,
  67. }) satisfies ProjectionDefinition<'cache-test/marks', MarksState>
  68. const secretUnit = {
  69. key: 'cache-test/secret',
  70. stateSchema: z.string(),
  71. init: () => '',
  72. apply: state => state,
  73. stateVersion: 1,
  74. } satisfies ProjectionDefinition<'cache-test/secret', string>
  75. /** One session's record document on the per-record medium. */
  76. const recordPath = (root: string, id: Session['id']): string =>
  77. join(root, projectionCacheDomainSpec.name, 'sessions', `${String(id)}.json`)
  78. /** Header shape for cachedSnapshot calls. */
  79. const headerOf = (id: SessionId, createdAt = 0, cwd?: string) =>
  80. ({ version: 0, id, createdAt, isSeeded: false, ...cwd === undefined ? {} : { cwd } })
  81. interface HarnessOptions {
  82. root?: string
  83. config?: { writeEveryEvents: number; writeIntervalMs: number }
  84. stateVersion?: number
  85. }
  86. const contexts: Context[] = []
  87. const roots: string[] = []
  88. async function harness(options: HarnessOptions = {}) {
  89. const root = options.root ?? await mkdtemp(join(tmpdir(), 'dsh-projcache-'))
  90. roots.push(root)
  91. const ctx = new Context()
  92. contexts.push(ctx)
  93. // The cache opens its domain through the storage stack; the json backend
  94. // lands the per-record tree under this tmp root.
  95. await ctx.plugin(Storage)
  96. await ctx.plugin({ name: storageJsonName, inject: storageJsonInject, apply: storageJsonApply, Config: storageJsonConfig }, { root })
  97. await ctx.plugin({ name: storageDomainName, inject: storageDomainInject, apply: storageDomainApply, Config: storageDomainConfig }, { backend: 'json' })
  98. await ctx.plugin(SessionStore)
  99. await ctx.plugin(SessionProjectionRegistry)
  100. ctx.sessionProjections.register(marksUnit(options.stateVersion))
  101. const fiber = await ctx.plugin(SessionProjectionCache, options.config ?? { writeEveryEvents: 100, writeIntervalMs: 60_000 })
  102. return { ctx, root, fiber, cache: ctx.sessionProjectionCache }
  103. }
  104. const mark = (session: Session, marks: string[]): SessionEvent =>
  105. session.append('cache-test/mark', { marks })
  106. const endTurn = (session: Session): SessionEvent =>
  107. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  108. /** The stored record for one session id (undefined = absent or unreadable). */
  109. async function storedRecord(root: string, id: Session['id']): Promise<CheckpointRecord | undefined> {
  110. try {
  111. const document = JSON.parse(await readFile(recordPath(root, id), 'utf8')) as { record: unknown }
  112. return checkpointRecord.parse(document.record)
  113. } catch {
  114. return undefined
  115. }
  116. }
  117. /** The stored rows for one session id (undefined = absent or unreadable). */
  118. async function storedRows(root: string, id: Session['id']): Promise<CheckpointRecord['rows'] | undefined> {
  119. return (await storedRecord(root, id))?.rows
  120. }
  121. /** Pre-seed one session's record document with a stored checkpoint record. */
  122. async function seedRecord(
  123. root: string,
  124. id: string,
  125. rows: CheckpointRecord['rows'],
  126. identity: CheckpointRecord['identity'] = {
  127. createdAt: 0,
  128. isSeeded: false,
  129. inheritedEventCount: SessionLogOffset(0),
  130. },
  131. ): Promise<void> {
  132. const path = recordPath(root, SessionId(id))
  133. await mkdir(dirname(path), { recursive: true })
  134. await writeFile(path, JSON.stringify({ version: projectionCacheDomainSpec.version, record: { identity, rows } }))
  135. }
  136. afterEach(async () => {
  137. vi.useRealTimers()
  138. await Promise.all(contexts.splice(0).map(ctx => ctx.fiber.dispose()))
  139. await Promise.all(roots.splice(0).map(root => rm(root, { recursive: true, force: true, maxRetries: 10, retryDelay: 100 })))
  140. })
  141. describe('SessionProjectionCache write policy', () => {
  142. it('writes a durable checkpoint at turn/end (mandatory point)', async () => {
  143. const { ctx, root } = await harness()
  144. const session = ctx.sessions.create(SessionId('turn-end'))
  145. mark(session, ['a'])
  146. // Creation already wrote the init cut; the mark is throttled, so the
  147. // stored row is still the creation-time cut (no marks folded).
  148. await vi.waitFor(async () => {
  149. expect((await storedRows(root, session.id))?.['cache-test/marks']?.seq).toBe(-1)
  150. }, { timeout: 5_000 })
  151. const end = endTurn(session)
  152. await vi.waitFor(async () => {
  153. expect((await storedRows(root, session.id))?.['cache-test/marks'])
  154. .toEqual({ ver: 1, seq: end.seq, val: { marks: ['a'] } })
  155. }, { timeout: 5_000 })
  156. })
  157. it('writes a checkpoint at session creation, capturing the seed-derived cut', async () => {
  158. const { ctx, root } = await harness()
  159. // A forked child seeded with its ancestor's title-like event: no
  160. // conversation follows, yet the creation write must capture the fold so
  161. // a crash or a live-held fork still lists the derived value.
  162. const session = ctx.sessions.create(SessionId('seeded'), {
  163. seed: [{ type: 'cache-test/mark', seq: 0, time: 1, data: { marks: ['seed'] } }] as SessionEvent[],
  164. })
  165. await vi.waitFor(async () => {
  166. expect((await storedRows(root, session.id))?.['cache-test/marks']?.val)
  167. .toEqual({ marks: ['seed'] })
  168. }, { timeout: 5_000 })
  169. })
  170. it('writes at session disposal (detach, the live-to-cold moment)', async () => {
  171. const { ctx, root } = await harness()
  172. // Sessions dispose with their owning fiber: create in a child plugin.
  173. let session: Session | undefined
  174. const owner = await ctx.plugin(Object.assign((inner: Context) => {
  175. session = inner.sessions.create(SessionId('detach'))
  176. }, { inject: ['sessions'] }))
  177. if (session === undefined) throw new Error('session was not created')
  178. mark(session, ['live'])
  179. await owner.dispose()
  180. const detached = session
  181. await vi.waitFor(async () => {
  182. expect((await storedRows(root, detached.id))?.['cache-test/marks']?.val).toEqual({ marks: ['live'] })
  183. }, { timeout: 5_000 })
  184. })
  185. it('flushes when the in-turn event count reaches the configured threshold', async () => {
  186. const { ctx, root } = await harness({ config: { writeEveryEvents: 3, writeIntervalMs: 60_000 } })
  187. const session = ctx.sessions.create(SessionId('count'))
  188. mark(session, ['1'])
  189. mark(session, ['2'])
  190. await vi.waitFor(async () => {
  191. expect((await storedRows(root, session.id))?.['cache-test/marks']?.seq).toBe(-1) // still the creation cut
  192. }, { timeout: 5_000 })
  193. mark(session, ['3'])
  194. await vi.waitFor(async () => {
  195. expect((await storedRows(root, session.id))?.['cache-test/marks']?.val).toEqual({ marks: ['3'] })
  196. }, { timeout: 5_000 })
  197. })
  198. it('flushes on the configured interval when the count threshold is not reached', async () => {
  199. const { ctx, cache } = await harness({ config: { writeEveryEvents: 100, writeIntervalMs: 20 } })
  200. const write = vi.spyOn(cache, 'write').mockResolvedValue()
  201. vi.useFakeTimers()
  202. const session = ctx.sessions.create(SessionId('interval'))
  203. write.mockClear()
  204. mark(session, ['slow'])
  205. await vi.advanceTimersByTimeAsync(19)
  206. expect(write).not.toHaveBeenCalled()
  207. await vi.advanceTimersByTimeAsync(1)
  208. expect(write).toHaveBeenCalledExactlyOnceWith(session)
  209. })
  210. it('write() on a never-dirty session checkpoints directly and rejects a non-JSON unit state', async () => {
  211. const { ctx, root } = await harness()
  212. // Never dirtied: no events — write() still lands the init-derived cut.
  213. const clean = ctx.sessions.create(SessionId('clean-write'))
  214. await ctx.sessionProjectionCache.write(clean)
  215. expect((await storedRows(root, clean.id))?.['cache-test/marks']).toEqual({ ver: 1, seq: -1, val: null })
  216. // A unit whose state violates the plain-JSON contract fails the write loud.
  217. ctx.sessionProjections.register({
  218. key: 'cache-test/marks2',
  219. stateSchema: z.custom<Map<string, string>>(() => true),
  220. init: () => new Map<string, string>(),
  221. apply: state => state,
  222. stateVersion: 1,
  223. })
  224. await expect(ctx.sessionProjectionCache.write(clean)).rejects.toThrow('not losslessly JSON-serializable')
  225. })
  226. it('plugin disposal clears armed interval timers and leaves cleaned sessions alone', async () => {
  227. vi.useFakeTimers()
  228. const { ctx, root, fiber } = await harness({ config: { writeEveryEvents: 100, writeIntervalMs: 5000 } })
  229. const armed = ctx.sessions.create(SessionId('armed'))
  230. const cleaned = ctx.sessions.create(SessionId('cleaned'))
  231. mark(armed, ['pending']) // timer armed, no write yet
  232. mark(cleaned, ['done'])
  233. endTurn(cleaned) // mandatory write; markClean leaves {pending: 0, timer: undefined} in the map
  234. await vi.advanceTimersByTimeAsync(0)
  235. await fiber.dispose()
  236. // The armed timer died with the plugin: advancing time writes nothing.
  237. await vi.advanceTimersByTimeAsync(10_000)
  238. // Only the creation cut exists: the armed mark never wrote.
  239. expect((await storedRows(root, armed.id))?.['cache-test/marks']?.seq).toBe(-1)
  240. })
  241. it('contains a durable write failure: logs a warning, event path unharmed, next write self-heals', async () => {
  242. const root = await mkdtemp(join(tmpdir(), 'dsh-projcache-'))
  243. roots.push(root)
  244. const ctx = new Context()
  245. contexts.push(ctx)
  246. await ctx.plugin(Storage)
  247. await ctx.plugin({ name: storageJsonName, inject: storageJsonInject, apply: storageJsonApply, Config: storageJsonConfig }, { root })
  248. await ctx.plugin({ name: storageDomainName, inject: storageDomainInject, apply: storageDomainApply, Config: storageDomainConfig }, { backend: 'json' })
  249. await ctx.plugin(SessionStore)
  250. await ctx.plugin(SessionProjectionRegistry)
  251. ctx.sessionProjections.register(marksUnit())
  252. await ctx.plugin(SessionProjectionCache, { writeEveryEvents: 100, writeIntervalMs: 60_000 })
  253. const warn = vi.spyOn(ctx.logger, 'warn').mockImplementation(() => {})
  254. // A directory where the record document must land makes the atomic
  255. // rename fail — including the creation write, so no row ever lands.
  256. const blocker = recordPath(root, SessionId('fail-soft'))
  257. await mkdir(blocker, { recursive: true })
  258. const session = ctx.sessions.create(SessionId('fail-soft'))
  259. mark(session, ['x'])
  260. endTurn(session)
  261. // The failed creation/turn-end writes are fire-and-forget: wait for the
  262. // warn (the write actually failed), then assert no row landed — the
  263. // property under test is that a failed write leaves no partial row.
  264. await vi.waitFor(() => {
  265. expect(warn).toHaveBeenCalledWith(expect.stringContaining('turn/end write for "fail-soft" failed'))
  266. }, { timeout: 5_000 })
  267. await vi.waitFor(async () => {
  268. expect(await storedRows(root, session.id)).toBeUndefined()
  269. }, { timeout: 5_000 })
  270. // Self-heal: once the blocker clears, the next mandatory point writes.
  271. await rm(recordPath(root, session.id), { recursive: true })
  272. mark(session, ['y'])
  273. endTurn(session)
  274. await vi.waitFor(async () => {
  275. expect((await storedRows(root, session.id))?.['cache-test/marks']?.val).toEqual({ marks: ['y'] })
  276. }, { timeout: 5_000 })
  277. })
  278. })
  279. describe('SessionProjectionCache listing read', () => {
  280. it('refuses a checkpoint created for a different inherited cut', async () => {
  281. const root = await mkdtemp(join(tmpdir(), 'dsh-projcache-'))
  282. roots.push(root)
  283. const id = SessionId('cut-identity')
  284. await seedRecord(
  285. root,
  286. id,
  287. { 'cache-test/marks': { ver: 1, seq: SessionSeq(1), val: { marks: ['seed'] } } },
  288. {
  289. createdAt: 0,
  290. isSeeded: true,
  291. inheritedEventCount: SessionLogOffset(2),
  292. },
  293. )
  294. const { cache } = await harness({ root })
  295. const seededHeader = { ...headerOf(id), isSeeded: true }
  296. expect(cache.cachedSnapshot(seededHeader, SessionLogOffset(2))?.values['cache-test/marks'])
  297. .toEqual({ marks: ['seed'] })
  298. expect(cache.cachedSnapshot(seededHeader, SessionLogOffset(1))).toBeUndefined()
  299. expect(() => cache.cachedSnapshot(headerOf(id), SessionLogOffset(1)))
  300. .toThrow('unseeded projection-cache identity inherited event count must be 0')
  301. })
  302. it('serves a creation-time checkpoint at the before-first-event cursor', async () => {
  303. const root = await mkdtemp(join(tmpdir(), 'dsh-projcache-'))
  304. roots.push(root)
  305. await seedRecord(root, 'before-first-event', {
  306. 'cache-test/marks': { ver: 1, seq: -1, val: null },
  307. })
  308. const { cache } = await harness({ root })
  309. expect(cache.cachedSnapshot(headerOf(SessionId('before-first-event')), SessionLogOffset(0)))
  310. .toEqual({ asOfSeq: -1, values: { 'cache-test/marks': { marks: [] } } })
  311. })
  312. it('keeps host-only checkpoint state out of cached wire snapshots', async () => {
  313. const root = await mkdtemp(join(tmpdir(), 'dsh-projcache-'))
  314. roots.push(root)
  315. await seedRecord(root, 'host-state', {
  316. 'cache-test/marks': { ver: 1, seq: SessionSeq(4), val: { marks: ['wire'] } },
  317. 'cache-test/secret': { ver: 1, seq: SessionSeq(4), val: 'private prompt text' },
  318. })
  319. const { ctx, cache } = await harness({ root })
  320. ctx.sessionProjections.register(secretUnit)
  321. const header = headerOf(SessionId('host-state'))
  322. expect(cache.cachedSnapshot(header, SessionLogOffset(0))).toEqual({
  323. asOfSeq: 4,
  324. values: { 'cache-test/marks': { marks: ['wire'] } },
  325. })
  326. expect(JSON.stringify(cache.cachedSnapshot(header, SessionLogOffset(0))))
  327. .not.toContain('private prompt text')
  328. })
  329. it('serves identity-matching rows with the cut watermark and refuses unrelated ones', async () => {
  330. const root = await mkdtemp(join(tmpdir(), 'dsh-projcache-'))
  331. roots.push(root)
  332. await seedRecord(root, 'listed', {
  333. 'cache-test/marks': { ver: 1, seq: SessionSeq(4), val: { marks: ['t'] } },
  334. })
  335. const { cache } = await harness({ root })
  336. const id = SessionId('listed')
  337. // Matching header: values plus the watermark the client seeds under.
  338. expect(cache.cachedSnapshot(headerOf(id), SessionLogOffset(0)))
  339. .toEqual({ asOfSeq: 4, values: { 'cache-test/marks': { marks: ['t'] } } })
  340. // A recreated id (different createdAt): the record is unrelated — no block.
  341. expect(cache.cachedSnapshot(headerOf(id, 777), SessionLogOffset(0))).toBeUndefined()
  342. // Unknown id: no block.
  343. expect(cache.cachedSnapshot(headerOf(SessionId('never-cached')), SessionLogOffset(0)))
  344. .toBeUndefined()
  345. })
  346. it('returns undefined when the stored record is version-mismatched', async () => {
  347. const root = await mkdtemp(join(tmpdir(), 'dsh-projcache-'))
  348. roots.push(root)
  349. // A stale version-stamped document is discarded at open: absent record.
  350. const path = recordPath(root, SessionId('all-stale'))
  351. await mkdir(dirname(path), { recursive: true })
  352. await writeFile(path, JSON.stringify({
  353. version: projectionCacheDomainSpec.version + 1,
  354. record: {
  355. identity: { createdAt: 0, isSeeded: false, inheritedEventCount: 0 },
  356. rows: { 'cache-test/marks': { ver: 1, seq: 4, val: { marks: ['old'] } } },
  357. },
  358. }))
  359. const { cache } = await harness({ root })
  360. expect(cache.cachedSnapshot(headerOf(SessionId('all-stale')), SessionLogOffset(0)))
  361. .toBeUndefined()
  362. })
  363. it('returns undefined when every stored row is version-mismatched', async () => {
  364. const root = await mkdtemp(join(tmpdir(), 'dsh-projcache-'))
  365. roots.push(root)
  366. // A current document whose rows all fail the live unit's stateVersion:
  367. // the listing view is empty, so no block is served.
  368. await seedRecord(root, 'row-stale', {
  369. 'cache-test/marks': { ver: 99, seq: SessionSeq(4), val: { marks: ['old'] } },
  370. })
  371. const { cache } = await harness({ root })
  372. expect(cache.cachedSnapshot(headerOf(SessionId('row-stale')), SessionLogOffset(0)))
  373. .toBeUndefined()
  374. })
  375. it('binds identity on cwd too: a matching cwd serves, a moved session does not', async () => {
  376. const root = await mkdtemp(join(tmpdir(), 'dsh-projcache-'))
  377. roots.push(root)
  378. await seedRecord(root, 'homed', {
  379. 'cache-test/marks': { ver: 1, seq: SessionSeq(2), val: { marks: ['w'] } },
  380. }, {
  381. createdAt: 0,
  382. cwd: '/work',
  383. isSeeded: false,
  384. inheritedEventCount: SessionLogOffset(0),
  385. })
  386. const { cache } = await harness({ root })
  387. const id = SessionId('homed')
  388. expect(cache.cachedSnapshot(headerOf(id, 0, '/work'), SessionLogOffset(0))?.values['cache-test/marks'])
  389. .toEqual({ marks: ['w'] })
  390. expect(cache.cachedSnapshot(headerOf(id, 0, '/elsewhere'), SessionLogOffset(0))).toBeUndefined()
  391. expect(cache.cachedSnapshot(headerOf(id, 0), SessionLogOffset(0))).toBeUndefined()
  392. })
  393. it('returns undefined for a malformed record document (refold from the log on the caller side)', async () => {
  394. const root = await mkdtemp(join(tmpdir(), 'dsh-projcache-'))
  395. roots.push(root)
  396. const path = recordPath(root, SessionId('malformed'))
  397. await mkdir(dirname(path), { recursive: true })
  398. await writeFile(path, 'not json at all')
  399. const { cache } = await harness({ root })
  400. expect(cache.cachedSnapshot(headerOf(SessionId('malformed')), SessionLogOffset(0)))
  401. .toBeUndefined()
  402. })
  403. })
  404. describe('SessionProjectionCache cold-read seeding', () => {
  405. /** One session's event log: turn/start, one mark per group, turn/end. */
  406. const storedLog = (marks: string[][]): SessionEvent[] => {
  407. const events: SessionEvent[] = [
  408. { type: 'turn/start', seq: SessionSeq(0), time: 0, data: { turn: 1 } },
  409. ]
  410. for (const m of marks) {
  411. events.push({
  412. type: 'cache-test/mark',
  413. seq: SessionSeq(events.length),
  414. time: events.length,
  415. data: { marks: m },
  416. })
  417. }
  418. events.push({
  419. type: 'turn/end',
  420. seq: SessionSeq(events.length),
  421. time: events.length,
  422. data: { turn: 1, reason: { kind: 'completed' } },
  423. })
  424. return events
  425. }
  426. it('hydratePrepared seeds from a matching row and retries from the exact log on a malformed one', async () => {
  427. const root = await mkdtemp(join(tmpdir(), 'dsh-projcache-'))
  428. roots.push(root)
  429. // Records land on disk before the domain opens, so the in-memory table
  430. // picks them up at init.
  431. await seedRecord(root, 'prepared-seeded', {
  432. 'cache-test/marks': { ver: 1, seq: SessionSeq(1), val: { marks: ['cached'] } },
  433. })
  434. await seedRecord(root, 'prepared-fallback', {
  435. 'cache-test/marks': { ver: 1, seq: SessionSeq(1), val: { marks: 'malformed' } },
  436. })
  437. const { cache } = await harness({ root })
  438. const events = storedLog([['fresh']])
  439. // A matching row hydrates the prepared Session without a persistence read.
  440. const seeded = headerOf(SessionId('prepared-seeded'))
  441. const seededSession = Session.create(seeded.id, events, seeded)
  442. expect(cache.hydratePrepared(seededSession, events)).toEqual({
  443. asOfSeq: 2,
  444. values: { 'cache-test/marks': { marks: ['cached'] } },
  445. })
  446. // A malformed row cannot seed the fold; hydration falls back to the
  447. // exact log so a valid Session stays readable.
  448. const fallback = headerOf(SessionId('prepared-fallback'))
  449. const fallbackSession = Session.create(fallback.id, events, fallback)
  450. expect(cache.hydratePrepared(fallbackSession, events)).toEqual({
  451. asOfSeq: 2,
  452. values: { 'cache-test/marks': { marks: ['fresh'] } },
  453. })
  454. // No row at all: hydrate from init over the exact log.
  455. const bare = headerOf(SessionId('prepared-bare'))
  456. const bareSession = Session.create(bare.id, events, bare)
  457. expect(cache.hydratePrepared(bareSession, events)).toEqual({
  458. asOfSeq: 2,
  459. values: { 'cache-test/marks': { marks: ['fresh'] } },
  460. })
  461. })
  462. it('coldSnapshot traverses the full log but applies only the events after each cached watermark', async () => {
  463. const root = await mkdtemp(join(tmpdir(), 'dsh-projcache-'))
  464. roots.push(root)
  465. // A cached row covering the prefix through seq 2 (three applies folded).
  466. await seedRecord(root, 'cold-snap', {
  467. 'cache-test/count': { ver: 1, seq: SessionSeq(2), val: 3 },
  468. }, {
  469. createdAt: 9,
  470. isSeeded: false,
  471. inheritedEventCount: SessionLogOffset(0),
  472. })
  473. const { cache, ctx } = await harness({ root })
  474. const apply = vi.fn((_state: number, _event: SessionEvent) => 1)
  475. ctx.sessionProjections.register({
  476. key: 'cache-test/count',
  477. stateSchema: z.number().int().nonnegative(),
  478. init: () => 0,
  479. apply,
  480. stateVersion: 1,
  481. } satisfies ProjectionDefinition<'cache-test/count', number>)
  482. const meta = headerOf(SessionId('cold-snap'), 9)
  483. const events = Array.from({ length: 5 }, (_, seq) => ({
  484. type: 'cache-test/mark', seq: SessionSeq(seq), time: seq, data: { marks: [`m${seq}`] },
  485. })) as SessionEvent[]
  486. const snapshot = cache.coldSnapshot(meta, SessionLogOffset(0), events)
  487. // The full log was traversed, but the fold applied only seqs 3 and 4.
  488. expect(apply).toHaveBeenCalledTimes(2)
  489. expect(apply.mock.calls.map(call => call[1].seq)).toEqual([3, 4])
  490. expect(snapshot.asOfSeq).toBe(4)
  491. // Host-only unit: folded but not served; the refreshed row is written
  492. // back (fail-soft, fire-and-forget) once the write lands.
  493. expect(Object.keys(snapshot.values)).not.toContain('cache-test/count')
  494. await vi.waitFor(async () => {
  495. expect((await storedRows(root, meta.id))?.['cache-test/count']?.seq).toBe(4)
  496. })
  497. // No cached row yet: the first cold read folds from init over the full
  498. // log and creates the cache row (the `?? {}` seed path).
  499. const fresh = headerOf(SessionId('cold-fresh'), 10)
  500. cache.coldSnapshot(fresh, SessionLogOffset(0), events)
  501. expect(apply).toHaveBeenCalledTimes(7) // 2 tail + 5 full
  502. await vi.waitFor(async () => {
  503. expect((await storedRows(root, fresh.id))?.['cache-test/count']?.seq).toBe(4)
  504. })
  505. })
  506. it('coldSnapshot write-back is fail-soft: a failed durable write logs and never throws', async () => {
  507. const root = await mkdtemp(join(tmpdir(), 'dsh-projcache-'))
  508. roots.push(root)
  509. const ctx = new Context()
  510. contexts.push(ctx)
  511. await ctx.plugin(Storage)
  512. await ctx.plugin({ name: storageJsonName, inject: storageJsonInject, apply: storageJsonApply, Config: storageJsonConfig }, { root })
  513. await ctx.plugin({ name: storageDomainName, inject: storageDomainInject, apply: storageDomainApply, Config: storageDomainConfig }, { backend: 'json' })
  514. await ctx.plugin(SessionStore)
  515. await ctx.plugin(SessionProjectionRegistry)
  516. ctx.sessionProjections.register(marksUnit())
  517. await ctx.plugin(SessionProjectionCache, { writeEveryEvents: 100, writeIntervalMs: 60_000 })
  518. const warn = vi.spyOn(ctx.logger, 'warn').mockImplementation(() => {})
  519. // A directory where the record document must land makes the write-back
  520. // fail; the cold read itself still succeeds and never throws.
  521. const meta = headerOf(SessionId('cold-fail'))
  522. await mkdir(recordPath(root, meta.id), { recursive: true })
  523. expect(ctx.sessionProjectionCache.coldSnapshot(meta, SessionLogOffset(0), [])).toBeDefined()
  524. // The failed write-back is fire-and-forget: poll for the warn instead of
  525. // assuming a fixed settle window (slow runners exceed it).
  526. await vi.waitFor(() => {
  527. expect(warn).toHaveBeenCalledWith(expect.stringContaining('cold-read write-back for "cold-fail" failed'))
  528. }, { timeout: 5_000 })
  529. })
  530. })