multi-edge-publication.spec.ts 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301
  1. /** Durable composition of historical chunk collapse and V3 system/reference migration. */
  2. import { Context } from '@deepseek-ai/cordis'
  3. import { Session, SessionId, SessionLogOffset, SessionSeq } from '@deepseek-ai/dsh-session'
  4. import type { SessionEvent, SessionHeader } from '@deepseek-ai/dsh-session'
  5. import { createSessionFormatCatalog } from '@deepseek-ai/dsh-session-format'
  6. import { releasedV0SessionFormatCodec, releasedV1SessionFormatCodec, sessionFormatV0ToV1 } from '@deepseek-ai/dsh-session-format-v0-to-v1'
  7. import {
  8. assertReleasedV2Header, RELEASED_V2_EVENT_TYPES, releasedV2SessionFormatCodec,
  9. restoreReleasedV2Artifact, sessionFormatV1ToV2,
  10. } from '@deepseek-ai/dsh-session-format-v1-to-v2'
  11. import { SessionFormatUnsupportedError } from '@deepseek-ai/dsh-session-persistence'
  12. import JsonlSessionPersistence from '@deepseek-ai/dsh-session-persistence-jsonl'
  13. import { appendFile, mkdir, mkdtemp, readFile, readdir, rm, stat, writeFile } from 'node:fs/promises'
  14. import { tmpdir } from 'node:os'
  15. import { basename, dirname, join } from 'node:path'
  16. import { scheduler } from 'node:timers/promises'
  17. import { afterEach, describe, expect, it, vi } from 'vitest'
  18. import { JsonlGenerationSourceChangedError } from '../src/generation.ts'
  19. import { generationLogPath, type JsonlCompression } from '../src/format.ts'
  20. import { compressZstdFrame, decompressZstdFrame, scanZstdFrames } from '../src/zstd.ts'
  21. const id = SessionId('multi-edge-seeded')
  22. const config = { provider: 'mock', model: 'mock' }
  23. const roots: string[] = []
  24. const contexts: Context[] = []
  25. const v2EventTypes = new Set(RELEASED_V2_EVENT_TYPES)
  26. const v2Catalog = createSessionFormatCatalog({
  27. currentVersion: 2,
  28. codecs: [releasedV0SessionFormatCodec, releasedV1SessionFormatCodec, releasedV2SessionFormatCodec],
  29. currentEncoder: releasedV2SessionFormatCodec,
  30. migrations: [sessionFormatV0ToV1, sessionFormatV1ToV2],
  31. restoreCurrent: artifact => restoreReleasedV2Artifact(artifact, v2EventTypes),
  32. restoreTransformedCurrent: artifact => restoreReleasedV2Artifact(artifact, v2EventTypes),
  33. restoreCurrentHeader(header) {
  34. assertReleasedV2Header(header)
  35. return header
  36. },
  37. })
  38. afterEach(async () => {
  39. try {
  40. await Promise.all(contexts.splice(0).map(ctx => ctx.fiber.dispose()))
  41. } finally {
  42. await Promise.all(roots.splice(0).map(root => rm(root, { recursive: true, force: true })))
  43. }
  44. })
  45. function message(role: 'user' | 'assistant', text: string) {
  46. return {
  47. id: text, role, content: [{ type: 'text', text }],
  48. source: role === 'user' ? { kind: 'user' } : { kind: 'model', ...config },
  49. }
  50. }
  51. function event(type: string, seq: number, data: object) {
  52. return { type, seq, time: 100 + seq, data }
  53. }
  54. function request(seq: number, system: string) {
  55. return event('request/header', seq, { header: { config, system }, reason: 'change' })
  56. }
  57. /** Packed rows consume four chunk coordinates before the inherited Assistant message. */
  58. function historicalRows() {
  59. return [
  60. event('turn/start', 0, { turn: 1 }),
  61. event('step/start', 1, { turn: 1, step: 1 }),
  62. { ...event('user/message', 2, message('user', 'question')), surfaceOp: 'append' },
  63. request(3, 'seed prompt'),
  64. { type: 'text-chunks', seq0: 4, time0: 104, data: { turn: 1, step: 1, index: 0, dt: [1, 1], texts: ['he', 'l', 'lo'] } },
  65. event('assistant/chunk', 7, { turn: 1, step: 1, chunk: { type: 'finish', reason: { kind: 'stop' } } }),
  66. { ...event('assistant/message', 8, { turn: 1, step: 1, message: message('assistant', 'hello') }), surfaceOp: 'append', sourceEventSeqs: [[4, 7]] },
  67. event('step/end', 9, { turn: 1, step: 1 }),
  68. event('turn/end', 10, { turn: 1, reason: { kind: 'completed' } }),
  69. event('session/end-seed', 11, {}),
  70. event('turn/start', 12, { turn: 2 }),
  71. event('step/start', 13, { turn: 2, step: 1 }),
  72. { ...event('user/message', 14, message('user', 'follow-up')), surfaceOp: 'append' },
  73. request(15, 'changed prompt'),
  74. event('compaction/prune', 16, { shadowedRange: { start: 2, end: 8 }, shadowedSeqs: [2, 8], shadowedTokenCount: 20 }),
  75. { ...event('user/message', 17, message('user', 'summary')), surfaceOp: { op: 'replace', start: 2, end: 8 }, sourceEventSeqs: [2, 8] },
  76. event('command/run', 18, { commandId: 'command', name: 'test', source: { kind: 'user' } }),
  77. event('command/done', 19, { commandId: 'command', kind: 'success', sourceEventSeq: 17 }),
  78. event('session/title', 20, { title: 'title', messageSeqs: [2, 14], source: { kind: 'fallback' } }),
  79. request(21, 'changed prompt'),
  80. event('step/end', 22, { turn: 2, step: 1 }),
  81. event('turn/end', 23, { turn: 2, reason: { kind: 'completed' } }),
  82. ]
  83. }
  84. async function mount(root: string, compression: JsonlCompression) {
  85. const ctx = new Context()
  86. contexts.push(ctx)
  87. await ctx.plugin(JsonlSessionPersistence, { root, compression })
  88. return ctx
  89. }
  90. async function seed(version: 0 | 1, compression: JsonlCompression, refuse = false) {
  91. const root = await mkdtemp(join(tmpdir(), 'dsh-multi-edge-publication-'))
  92. roots.push(root)
  93. const path = generationLogPath(root, undefined, id, version, compression)
  94. await mkdir(dirname(path), { recursive: true })
  95. const header = { type: 'session', version, id, createdAt: 1, parentSession: 'parent', delegationDepth: 0, seedLength: 11 }
  96. const rows = refuse ? [...historicalRows().slice(0, -1), request(23, 'outside step')] : historicalRows()
  97. const headerLine = JSON.stringify(header) + '\n'
  98. const body = rows.map(row => JSON.stringify(row)).join('\n') + '\n'
  99. const bytes = compression === 'none' ? Buffer.from(headerLine + body) : Buffer.concat([
  100. await compressZstdFrame(headerLine), await compressZstdFrame(body),
  101. ])
  102. await writeFile(path, bytes)
  103. return { root, path, header, rows }
  104. }
  105. async function observe(path: string) {
  106. const identity = await stat(path, { bigint: true })
  107. return {
  108. bytes: await readFile(path), dev: identity.dev, ino: identity.ino,
  109. size: identity.size, mtimeNs: identity.mtimeNs, ctimeNs: identity.ctimeNs,
  110. }
  111. }
  112. async function readSession(ctx: Context, access: 'read' | 'write') {
  113. const handle = await ctx.sessionPersistence.open(id, access)
  114. try {
  115. const result = await handle.read()
  116. const session = Session.fromRestore(id, result.events, handle.header, handle.inheritedEventCount, result.eventState)
  117. if (access === 'write') await handle.flush()
  118. return { header: handle.header, events: result.events, cut: handle.inheritedEventCount, session }
  119. } finally {
  120. await handle.close()
  121. }
  122. }
  123. function visible(role: 'system' | 'user' | 'assistant', text: string) {
  124. return { role, content: [{ type: 'text', text }] }
  125. }
  126. function assertRequests(events: readonly SessionEvent[], header: SessionHeader) {
  127. const requests = events.filter(event => event.type === 'request/header')
  128. expect(requests.map(event => event.data.header)).toEqual([{ config }, { config }, { config }])
  129. expect(requests.map(event => Session.fromRestore(
  130. id, events.slice(0, event.seq + 1), header, SessionLogOffset(0), 'shared-frozen',
  131. ).deriveMessages().map(({ role, content }) => ({ role, content })))).toEqual([
  132. [visible('system', 'seed prompt'), visible('user', 'question')],
  133. [visible('system', 'changed prompt'), visible('user', 'question'), visible('assistant', 'hello'), visible('user', 'follow-up')],
  134. [visible('system', 'changed prompt'), visible('user', 'summary'), visible('user', 'follow-up')],
  135. ])
  136. }
  137. function assertMigrated(result: Awaited<ReturnType<typeof readSession>>) {
  138. const { events, header, cut, session } = result
  139. expect(header).toEqual({ version: 3, id, createdAt: 1, parentSession: 'parent', delegationDepth: 0, isSeeded: true })
  140. expect(events.map(event => event.seq)).toEqual(Array.from({ length: 23 }, (_, seq) => seq))
  141. expect(events.filter(event => event.type.startsWith('assistant/'))).toEqual([{
  142. type: 'assistant/message', seq: 6, time: 108, surfaceOp: 'append',
  143. data: { turn: 1, step: 1, message: message('assistant', 'hello'), stream: [
  144. { type: 'text-chunks', time0: 104, index: 0, dt: [1, 1], texts: ['he', 'l', 'lo'] },
  145. { type: 'chunk', time: 107, chunk: { type: 'finish', reason: { kind: 'stop' } } },
  146. ] },
  147. }])
  148. expect(events.filter(event => event.type === 'system/message').map(event => ({
  149. seq: event.seq, content: event.data.message.content, surfaceOp: event.surfaceOp, sourceEventSeqs: event.sourceEventSeqs,
  150. }))).toEqual([
  151. { seq: 2, content: [], surfaceOp: 'append', sourceEventSeqs: undefined },
  152. { seq: 4, content: visible('system', 'seed prompt').content, surfaceOp: { op: 'replace', startSeq: 2, endSeq: 2 }, sourceEventSeqs: [2] },
  153. { seq: 13, content: visible('system', 'changed prompt').content, surfaceOp: { op: 'replace', startSeq: 4, endSeq: 4 }, sourceEventSeqs: [4] },
  154. ])
  155. expect(events[15]).toMatchObject({ type: 'compaction/prune', data: { shadowedRange: { start: 3, end: 6 }, shadowedSeqs: [3, 6] } })
  156. expect(events[16]).toMatchObject({ type: 'user/message', surfaceOp: { op: 'replace', startSeq: 3, endSeq: 6 }, sourceEventSeqs: [3, 6] })
  157. expect(events[18]).toMatchObject({ type: 'command/done', data: { sourceEventSeq: 16 } })
  158. expect(events[19]).toMatchObject({ type: 'session/title', data: { messageSeqs: [3, 12] } })
  159. expect(cut).toBe(9)
  160. expect(events[9]).toEqual({ type: 'session/end-seed', seq: 9, time: 111, data: { inherited: true } })
  161. expect(session.inheritedEventCount).toBe(9)
  162. expect(session.firstLiveSeq).toBe(23)
  163. expect(session.isOwnSeq(SessionSeq(8))).toBe(false)
  164. expect(session.isOwnSeq(SessionSeq(9))).toBe(true)
  165. expect(session.ownEvents()).toEqual([
  166. ...events.slice(9),
  167. expect.objectContaining({ type: 'session/end-seed', seq: 23, data: {} }),
  168. ])
  169. expect(session.surface.nodes).toEqual([13, 16, 12])
  170. expect(session.deriveMessages().map(({ role, content }) => ({ role, content }))).toEqual([
  171. visible('system', 'changed prompt'), visible('user', 'summary'), visible('user', 'follow-up'),
  172. ])
  173. assertRequests(events, header)
  174. }
  175. async function publishedRows(path: string, compression: JsonlCompression) {
  176. const bytes = await readFile(path)
  177. let plaintext = bytes
  178. if (compression === 'zstd') {
  179. const { frames, tornStart } = scanZstdFrames(bytes)
  180. expect(tornStart).toBeUndefined()
  181. expect(frames.length).toBeGreaterThan(0)
  182. plaintext = Buffer.concat(await Promise.all(frames.map(frame => decompressZstdFrame(bytes.subarray(frame.start, frame.end)))))
  183. }
  184. return plaintext.toString('utf8').trimEnd().split('\n').map((line): unknown => JSON.parse(line))
  185. }
  186. describe.each([0, 1] as const)('V%s multi-edge durable publication', (version) => {
  187. it.each(['none', 'zstd'] as const)('publishes only V3 after chunk collapse, system changes, and reference remapping (%s)', async (compression) => {
  188. const { root, path } = await seed(version, compression)
  189. const source = await observe(path)
  190. const ctx = await mount(root, compression)
  191. const prepared = await readSession(ctx, 'read')
  192. assertMigrated(prepared)
  193. await ctx.sessionPersistence.flush()
  194. expect(await observe(path)).toEqual(source)
  195. expect(await readdir(dirname(path))).toEqual([basename(path)])
  196. const written = await readSession(ctx, 'write')
  197. assertMigrated(written)
  198. expect(written.events).toEqual(prepared.events)
  199. await ctx.fiber.dispose()
  200. contexts.splice(contexts.indexOf(ctx), 1)
  201. const successor = generationLogPath(root, undefined, id, 3, compression)
  202. expect((await readdir(dirname(path))).filter(name => name !== 'session.lock').sort())
  203. .toEqual([basename(path), basename(successor)].sort())
  204. expect(await publishedRows(successor, compression)).toEqual([{ type: 'session', ...prepared.header }, ...prepared.events])
  205. expect(await observe(path)).toEqual(source)
  206. const published = await observe(successor)
  207. const reopened = await mount(root, compression)
  208. const native = await readSession(reopened, 'read')
  209. assertMigrated(native)
  210. expect(native.events).toEqual(prepared.events)
  211. expect(native.cut).toBe(prepared.cut)
  212. const repeated = await readSession(reopened, 'write')
  213. expect(repeated.events).toEqual(prepared.events)
  214. expect(repeated.cut).toBe(prepared.cut)
  215. await reopened.sessionPersistence.flush()
  216. expect(await observe(path)).toEqual(source)
  217. expect(await observe(successor)).toEqual(published)
  218. })
  219. it.each(['none', 'zstd'] as const)('re-prepares populated history after source drift rejects stale publication (%s)', async (compression) => {
  220. const { root, path } = await seed(version, compression)
  221. const source = await observe(path)
  222. const ctx = await mount(root, compression)
  223. const prepared = await readSession(ctx, 'read')
  224. assertMigrated(prepared)
  225. const tail = event('feedback/record', 24, { text: 'arrived after preparation' })
  226. const line = JSON.stringify(tail) + '\n'
  227. const appended = compression === 'none' ? Buffer.from(line) : await compressZstdFrame(line)
  228. const yieldSpy = vi.spyOn(scheduler, 'yield').mockImplementationOnce(async () => {
  229. await appendFile(path, appended)
  230. })
  231. try {
  232. await expect(readSession(ctx, 'write')).rejects.toBeInstanceOf(JsonlGenerationSourceChangedError)
  233. expect(yieldSpy).toHaveBeenCalled()
  234. } finally {
  235. yieldSpy.mockRestore()
  236. }
  237. const changed = await observe(path)
  238. expect(changed.bytes).toEqual(Buffer.concat([source.bytes, appended]))
  239. expect(changed).toMatchObject({ dev: source.dev, ino: source.ino })
  240. expect((await readdir(dirname(path))).filter(name => name !== 'session.lock')).toEqual([basename(path)])
  241. const retried = await readSession(ctx, 'write')
  242. const expected = [...prepared.events, { ...tail, seq: 23 }]
  243. expect(retried.events).toEqual(expected)
  244. expect(retried.cut).toBe(prepared.cut)
  245. assertRequests(retried.events, retried.header)
  246. const successor = generationLogPath(root, undefined, id, 3, compression)
  247. expect(await publishedRows(successor, compression)).toEqual([{ type: 'session', ...prepared.header }, ...expected])
  248. expect((await readdir(dirname(path))).filter(name => name !== 'session.lock').sort())
  249. .toEqual([basename(path), basename(successor)].sort())
  250. await ctx.fiber.dispose()
  251. contexts.splice(contexts.indexOf(ctx), 1)
  252. const reopened = await mount(root, compression)
  253. const native = await readSession(reopened, 'read')
  254. expect(native.events).toEqual(expected)
  255. expect(native.cut).toBe(prepared.cut)
  256. assertRequests(native.events, native.header)
  257. expect(await observe(path)).toEqual(changed)
  258. })
  259. it.each(['none', 'zstd'] as const)('refuses a late V3 prompt outside a step without publishing earlier edges (%s)', async (compression) => {
  260. const { root, path, header, rows } = await seed(version, compression, true)
  261. const restoreV2 = v2Catalog.createRestore(header, { recovery: 'strict', validation: 'current' })
  262. for (const row of rows) restoreV2.decodeRow(row)
  263. const validV2 = restoreV2.finish()
  264. expect(validV2.inheritedEventCount).toBe(7)
  265. expect(validV2.events.slice(-2)).toEqual([
  266. { ...event('step/end', 22, { turn: 2, step: 1 }), seq: 18 },
  267. { ...request(23, 'outside step'), seq: 19 },
  268. ])
  269. const source = await observe(path)
  270. const ctx = await mount(root, compression)
  271. for (const access of ['read', 'write', 'read', 'write'] as const) {
  272. const failure = readSession(ctx, access)
  273. await expect(failure).rejects.toBeInstanceOf(SessionFormatUnsupportedError)
  274. await expect(failure).rejects.toThrow(/outside an open step/)
  275. await ctx.sessionPersistence.flush()
  276. expect(await observe(path)).toEqual(source)
  277. expect((await readdir(dirname(path))).filter(name => name !== 'session.lock')).toEqual([basename(path)])
  278. }
  279. })
  280. })