cache.spec.ts 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280
  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. * cached listing read. The durable medium is one `projection_cache.json`
  6. * per session under the cache's own configured root
  7. * (`<root>/<session-id>/projection_cache.json`); the cache never consults
  8. * the persistence layer.
  9. */
  10. import { afterEach, describe, expect, it, vi } from 'vitest'
  11. import { mkdtemp, mkdir, readFile, rm, writeFile } from 'node:fs/promises'
  12. import { tmpdir } from 'node:os'
  13. import { join } from 'node:path'
  14. import { Context } from '@deepseek-ai/cordis'
  15. import { z } from 'zod'
  16. import SessionStore, { SessionId } from '@deepseek-ai/dsh-session'
  17. import type { Session, SessionEvent } from '@deepseek-ai/dsh-session'
  18. import SessionProjectionRegistry from '@deepseek-ai/dsh-session-projection'
  19. import type { ProjectionDefinition } from '@deepseek-ai/dsh-session-projection'
  20. import SessionProjectionCache from '../src/index.ts'
  21. import { checkpointRecord } from '../src/spec.ts'
  22. import type { CheckpointRecord } from '../src/spec.ts'
  23. declare module '@deepseek-ai/dsh-session-projection/types' {
  24. interface SessionProjectionStateMap {
  25. 'cache-test/marks': MarksState
  26. 'cache-test/marks2': Map<string, string>
  27. }
  28. interface SessionProjectionMap {
  29. 'cache-test/marks': { marks: string[] }
  30. }
  31. }
  32. declare module '@deepseek-ai/dsh-session/types' {
  33. interface SessionEventMap {
  34. 'cache-test/mark': { marks: string[] }
  35. }
  36. interface OutOfBandSessionEventMap {
  37. 'cache-test/mark': true
  38. }
  39. }
  40. type MarksState = { marks: string[] } | null
  41. const marksUnit = (stateVersion = 1) => ({
  42. key: 'cache-test/marks',
  43. stateSchema: z.object({ marks: z.array(z.string()) }).nullable(),
  44. init: () => null,
  45. apply: (state, event) => (event.type === 'cache-test/mark' ? (event).data : state),
  46. wire: {
  47. viewSchema: z.object({ marks: z.array(z.string()) }),
  48. view: state => state ?? { marks: [] },
  49. },
  50. stateVersion,
  51. }) satisfies ProjectionDefinition<'cache-test/marks', MarksState>
  52. /** One session's cache file under the cache's own root. */
  53. const cachePath = (root: string, id: Session['id']): string =>
  54. join(root, String(id), 'projection_cache.json')
  55. /** Header shape for cachedSnapshot calls. */
  56. const headerOf = (id: SessionId, createdAt = 0, cwd?: string) =>
  57. ({ version: 0, id, createdAt, ...cwd === undefined ? {} : { cwd } })
  58. interface HarnessOptions {
  59. root?: string
  60. config?: { writeEveryEvents: number; writeIntervalMs: number }
  61. stateVersion?: number
  62. }
  63. const contexts: Context[] = []
  64. const roots: string[] = []
  65. async function harness(options: HarnessOptions = {}) {
  66. const root = options.root ?? await mkdtemp(join(tmpdir(), 'dsh-projcache-'))
  67. roots.push(root)
  68. const ctx = new Context()
  69. contexts.push(ctx)
  70. await ctx.plugin(SessionStore)
  71. await ctx.plugin(SessionProjectionRegistry)
  72. ctx.sessionProjections.register(marksUnit(options.stateVersion))
  73. const fiber = await ctx.plugin(SessionProjectionCache, {
  74. root,
  75. ...options.config ?? { writeEveryEvents: 100, writeIntervalMs: 60_000 },
  76. })
  77. return { ctx, root, fiber, cache: ctx.sessionProjectionCache }
  78. }
  79. const mark = (session: Session, marks: string[]): SessionEvent =>
  80. session.append('cache-test/mark', { marks })
  81. const endTurn = (session: Session): SessionEvent =>
  82. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  83. /** The stored record for one session id (undefined = absent or unreadable). */
  84. async function storedRecord(root: string, id: Session['id']): Promise<CheckpointRecord | undefined> {
  85. try {
  86. return checkpointRecord.parse(JSON.parse(await readFile(cachePath(root, id), 'utf8')))
  87. } catch {
  88. return undefined
  89. }
  90. }
  91. /** The stored rows for one session id (undefined = absent or unreadable). */
  92. async function storedRows(root: string, id: Session['id']): Promise<CheckpointRecord['rows'] | undefined> {
  93. return (await storedRecord(root, id))?.rows
  94. }
  95. /** Pre-seed one session's cache file with a stored checkpoint record. */
  96. async function seedRecord(
  97. root: string,
  98. id: string,
  99. rows: CheckpointRecord['rows'],
  100. identity: CheckpointRecord['identity'] = { createdAt: 0 },
  101. ): Promise<void> {
  102. await mkdir(join(root, id), { recursive: true })
  103. await writeFile(cachePath(root, SessionId(id)), JSON.stringify({ identity, rows }))
  104. }
  105. /** Wait until queued fail-soft writes (event-listener fire-and-forget over real fs I/O) drain. */
  106. const settle = () => new Promise(resolve => setTimeout(resolve, 40))
  107. afterEach(async () => {
  108. vi.useRealTimers()
  109. await Promise.all(contexts.splice(0).map(ctx => ctx.fiber.dispose()))
  110. await Promise.all(roots.splice(0).map(root => rm(root, { recursive: true, force: true })))
  111. })
  112. describe('SessionProjectionCache write policy', () => {
  113. it('writes a durable checkpoint at turn/end (mandatory point)', async () => {
  114. const { ctx, root } = await harness()
  115. const session = ctx.sessions.create(SessionId('turn-end'))
  116. mark(session, ['a'])
  117. expect(await storedRows(root, session.id)).toBeUndefined() // throttled: no write yet
  118. const end = endTurn(session)
  119. await settle()
  120. const rows = await storedRows(root, session.id)
  121. expect(rows?.['cache-test/marks']).toEqual({ ver: 1, seq: end.seq, val: { marks: ['a'] } })
  122. })
  123. it('writes at session disposal (detach, the live-to-cold moment)', async () => {
  124. const { ctx, root } = await harness()
  125. // Sessions dispose with their owning fiber: create in a child plugin.
  126. let session: Session | undefined
  127. const owner = await ctx.plugin(Object.assign((inner: Context) => {
  128. session = inner.sessions.create(SessionId('detach'))
  129. }, { inject: ['sessions'] }))
  130. if (session === undefined) throw new Error('session was not created')
  131. mark(session, ['live'])
  132. await owner.dispose()
  133. await settle()
  134. expect((await storedRows(root, session.id))?.['cache-test/marks']?.val).toEqual({ marks: ['live'] })
  135. })
  136. it('flushes when the in-turn event count reaches the configured threshold', async () => {
  137. const { ctx, root } = await harness({ config: { writeEveryEvents: 3, writeIntervalMs: 60_000 } })
  138. const session = ctx.sessions.create(SessionId('count'))
  139. mark(session, ['1'])
  140. mark(session, ['2'])
  141. await settle()
  142. expect(await storedRows(root, session.id)).toBeUndefined()
  143. mark(session, ['3'])
  144. await settle()
  145. expect((await storedRows(root, session.id))?.['cache-test/marks']?.val).toEqual({ marks: ['3'] })
  146. })
  147. it('flushes on the configured interval when the count threshold is not reached', async () => {
  148. const { ctx, root } = await harness({ config: { writeEveryEvents: 100, writeIntervalMs: 20 } })
  149. const session = ctx.sessions.create(SessionId('interval'))
  150. mark(session, ['slow'])
  151. await new Promise(resolve => setTimeout(resolve, 10)) // before the interval
  152. expect(await storedRows(root, session.id)).toBeUndefined()
  153. await settle() // past the interval; the fire-and-forget write lands
  154. expect((await storedRows(root, session.id))?.['cache-test/marks']?.val).toEqual({ marks: ['slow'] })
  155. })
  156. it('write() on a never-dirty session checkpoints directly and rejects a non-JSON unit state', async () => {
  157. const { ctx, root } = await harness()
  158. // Never dirtied: no events — write() still lands the init-derived cut.
  159. const clean = ctx.sessions.create(SessionId('clean-write'))
  160. await ctx.sessionProjectionCache.write(clean)
  161. expect((await storedRows(root, clean.id))?.['cache-test/marks']).toEqual({ ver: 1, seq: -1, val: null })
  162. // A unit whose state violates the plain-JSON contract fails the write loud.
  163. ctx.sessionProjections.register({
  164. key: 'cache-test/marks2',
  165. stateSchema: z.custom<Map<string, string>>(() => true),
  166. init: () => new Map<string, string>(),
  167. apply: state => state,
  168. stateVersion: 1,
  169. })
  170. await expect(ctx.sessionProjectionCache.write(clean)).rejects.toThrow('not losslessly JSON-serializable')
  171. })
  172. it('plugin disposal clears armed interval timers and leaves cleaned sessions alone', async () => {
  173. vi.useFakeTimers()
  174. const { ctx, root, fiber } = await harness({ config: { writeEveryEvents: 100, writeIntervalMs: 5000 } })
  175. const armed = ctx.sessions.create(SessionId('armed'))
  176. const cleaned = ctx.sessions.create(SessionId('cleaned'))
  177. mark(armed, ['pending']) // timer armed, no write yet
  178. mark(cleaned, ['done'])
  179. endTurn(cleaned) // mandatory write; markClean leaves {pending: 0, timer: undefined} in the map
  180. await vi.advanceTimersByTimeAsync(0)
  181. await fiber.dispose()
  182. // The armed timer died with the plugin: advancing time writes nothing.
  183. await vi.advanceTimersByTimeAsync(10_000)
  184. expect(await storedRows(root, armed.id)).toBeUndefined()
  185. })
  186. it('contains a durable write failure: logs a warning, event path unharmed, next write self-heals', async () => {
  187. const root = await mkdtemp(join(tmpdir(), 'dsh-projcache-'))
  188. roots.push(root)
  189. const ctx = new Context()
  190. contexts.push(ctx)
  191. await ctx.plugin(SessionStore)
  192. await ctx.plugin(SessionProjectionRegistry)
  193. ctx.sessionProjections.register(marksUnit())
  194. await ctx.plugin(SessionProjectionCache, { root, writeEveryEvents: 100, writeIntervalMs: 60_000 })
  195. const warn = vi.spyOn(ctx.logger, 'warn').mockImplementation(() => {})
  196. const session = ctx.sessions.create(SessionId('fail-soft'))
  197. // A directory where the cache file must land makes the atomic rename
  198. // fail on the first write...
  199. await mkdir(cachePath(root, session.id), { recursive: true })
  200. mark(session, ['x'])
  201. endTurn(session)
  202. await settle()
  203. expect(await storedRows(root, session.id)).toBeUndefined()
  204. expect(warn).toHaveBeenCalledWith(expect.stringContaining('turn/end write for "fail-soft" failed'))
  205. // Self-heal: once the blocker clears, the next mandatory point writes.
  206. await rm(cachePath(root, session.id), { recursive: true })
  207. mark(session, ['y'])
  208. endTurn(session)
  209. await settle()
  210. expect((await storedRows(root, session.id))?.['cache-test/marks']?.val).toEqual({ marks: ['y'] })
  211. })
  212. })
  213. describe('SessionProjectionCache listing read', () => {
  214. it('serves identity-matching rows with the cut watermark and refuses unrelated ones', async () => {
  215. const root = await mkdtemp(join(tmpdir(), 'dsh-projcache-'))
  216. roots.push(root)
  217. await seedRecord(root, 'listed', { 'cache-test/marks': { ver: 1, seq: 4, val: { marks: ['t'] } } })
  218. const { cache } = await harness({ root })
  219. const id = SessionId('listed')
  220. // Matching header: values plus the watermark the client seeds under.
  221. expect(await cache.cachedSnapshot(headerOf(id))).toEqual({ asOfSeq: 4, values: { 'cache-test/marks': { marks: ['t'] } } })
  222. // A recreated id (different createdAt): the record is unrelated — no block.
  223. expect(await cache.cachedSnapshot(headerOf(id, 777))).toBeUndefined()
  224. // Unknown id: no block.
  225. expect(await cache.cachedSnapshot(headerOf(SessionId('never-cached')))).toBeUndefined()
  226. })
  227. it('returns undefined when every stored row is version-mismatched', async () => {
  228. const root = await mkdtemp(join(tmpdir(), 'dsh-projcache-'))
  229. roots.push(root)
  230. await seedRecord(root, 'all-stale', { 'cache-test/marks': { ver: 99, seq: 4, val: { marks: ['old'] } } })
  231. const { cache } = await harness({ root })
  232. expect(await cache.cachedSnapshot(headerOf(SessionId('all-stale')))).toBeUndefined()
  233. })
  234. it('binds identity on cwd too: a matching cwd serves, a moved session does not', async () => {
  235. const root = await mkdtemp(join(tmpdir(), 'dsh-projcache-'))
  236. roots.push(root)
  237. await seedRecord(root, 'homed', { 'cache-test/marks': { ver: 1, seq: 2, val: { marks: ['w'] } } }, { createdAt: 0, cwd: '/work' })
  238. const { cache } = await harness({ root })
  239. const id = SessionId('homed')
  240. expect((await cache.cachedSnapshot(headerOf(id, 0, '/work')))?.values['cache-test/marks']).toEqual({ marks: ['w'] })
  241. expect(await cache.cachedSnapshot(headerOf(id, 0, '/elsewhere'))).toBeUndefined()
  242. expect(await cache.cachedSnapshot(headerOf(id, 0))).toBeUndefined()
  243. })
  244. it('returns undefined for a malformed cache file (refold from the log on the caller side)', async () => {
  245. const root = await mkdtemp(join(tmpdir(), 'dsh-projcache-'))
  246. roots.push(root)
  247. await mkdir(join(root, 'malformed'), { recursive: true })
  248. await writeFile(cachePath(root, SessionId('malformed')), 'not json at all')
  249. const { cache } = await harness({ root })
  250. expect(await cache.cachedSnapshot(headerOf(SessionId('malformed')))).toBeUndefined()
  251. })
  252. })