format-decoder.spec.ts 41 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101
  1. import { afterEach, describe, expect, it, vi } from 'vitest'
  2. import { SessionId } from '@deepseek-ai/dsh-session'
  3. import type { SessionEvent } from '@deepseek-ai/dsh-session'
  4. import {
  5. SessionPersistenceRevision,
  6. SessionPersistenceRevisionConflictError,
  7. } from '../src/revision.ts'
  8. import type {
  9. SessionFormatMigration,
  10. StoredEventReadCompletion,
  11. StoredSessionSource,
  12. } from '../src/format-decoder.ts'
  13. import { sessionFormatVersionRefusal } from '../src/format-decoder.ts'
  14. import { unversionedFormatCompatibility } from '../src/format-v0-compat.ts'
  15. const id = SessionId('format-migration')
  16. type SessionFormatMigrationInstance = InstanceType<SessionFormatMigration>
  17. function eventLog(): SessionEvent[] {
  18. return [
  19. { type: 'turn/start', seq: 0, time: 1, data: { turn: 1 } },
  20. { type: 'turn/end', seq: 1, time: 2, data: { turn: 1, reason: { kind: 'completed' } } },
  21. ]
  22. }
  23. async function collectEvents(events: AsyncIterable<SessionEvent>): Promise<SessionEvent[]> {
  24. const collected: SessionEvent[] = []
  25. for await (const event of events) collected.push(event)
  26. return collected
  27. }
  28. async function decodedFailure(
  29. decoded: ReturnType<typeof import('../src/format-decoder.ts')['decodeStoredSession']>,
  30. ): Promise<Error> {
  31. const completion = decoded.completed.catch((error: unknown) => error)
  32. const consumption = collectEvents(decoded.events).catch((error: unknown) => error)
  33. const [streamFailure, completionFailure] = await Promise.all([consumption, completion])
  34. expect(streamFailure).toBe(completionFailure)
  35. expect(streamFailure).toBeInstanceOf(Error)
  36. return streamFailure as Error
  37. }
  38. function storedSource(
  39. version: number,
  40. events: readonly unknown[],
  41. ): { source: StoredSessionSource<never>; reads: number[]; meta: Record<string, unknown> } {
  42. const reads: number[] = []
  43. const meta: Record<string, unknown> = { version, id, createdAt: 1 }
  44. return {
  45. meta,
  46. reads,
  47. source: {
  48. meta,
  49. revision: SessionPersistenceRevision(`format-v${version}`),
  50. readEvents({ fromSeq = 0 } = {}) {
  51. reads.push(fromSeq)
  52. return {
  53. events: (async function* (): AsyncIterable<unknown> {
  54. for (const event of events) {
  55. const seq = typeof event === 'object' && event !== null
  56. ? (event as { seq?: unknown }).seq
  57. : undefined
  58. if (!Number.isSafeInteger(seq) || (seq as number) < 0 || (seq as number) >= fromSeq) {
  59. yield structuredClone(event)
  60. }
  61. }
  62. })(),
  63. completed: Promise.resolve({}),
  64. }
  65. },
  66. },
  67. }
  68. }
  69. function defineMigration(
  70. from: number,
  71. create: () => SessionFormatMigrationInstance,
  72. to = from + 1,
  73. ): SessionFormatMigration {
  74. return class implements SessionFormatMigrationInstance {
  75. static readonly from = from
  76. static readonly to = to
  77. private readonly delegate = create()
  78. header(meta: unknown): unknown {
  79. return this.delegate.header(meta)
  80. }
  81. event(value: unknown): unknown {
  82. return this.delegate.event(value)
  83. }
  84. finish(): void {
  85. this.delegate.finish?.()
  86. }
  87. }
  88. }
  89. function migration(
  90. from: number,
  91. calls: string[],
  92. to = from + 1,
  93. ): SessionFormatMigration {
  94. return defineMigration(from, () => {
  95. let observedInput = false
  96. return {
  97. header(meta) {
  98. calls.push(`header:${from}`)
  99. return { ...(meta as Record<string, unknown>), version: to }
  100. },
  101. event(value) {
  102. if (!observedInput) {
  103. calls.push(`events:${from}`)
  104. observedInput = true
  105. }
  106. const event = value as SessionEvent
  107. const data = event.data as Record<string, unknown>
  108. const migrationPath = Array.isArray(data['migrationPath'])
  109. ? data['migrationPath'] as unknown[]
  110. : []
  111. return {
  112. ...event,
  113. data: {
  114. ...data,
  115. [`migratedFrom${from}`]: true,
  116. migrationPath: [...migrationPath, from],
  117. },
  118. }
  119. },
  120. }
  121. }, to)
  122. }
  123. async function configuredDecoder(
  124. currentVersion: number,
  125. migrations: readonly SessionFormatMigration[],
  126. calls: string[] = [],
  127. ): Promise<{
  128. decodeStoredSession: typeof import('../src/format-decoder.ts')['decodeStoredSession']
  129. decodeStoredSessionHeader: typeof import('../src/format-decoder.ts')['decodeStoredSessionHeader']
  130. validateHeader: ReturnType<typeof vi.fn>
  131. }> {
  132. vi.resetModules()
  133. const validateHeader = vi.fn((sessionId: SessionId, _seed: unknown, meta: unknown) => {
  134. calls.push('validate-header')
  135. const record = meta as Record<string, unknown>
  136. if (record['version'] !== currentVersion) {
  137. throw new Error(`current header validator received v${String(record['version'])}`)
  138. }
  139. if (record['id'] !== sessionId) throw new Error('current header validator received the wrong id')
  140. if (!Number.isSafeInteger(record['createdAt'])) {
  141. throw new Error('current header validator received invalid createdAt')
  142. }
  143. return { header: Object.freeze(structuredClone(record)) }
  144. })
  145. vi.doMock('@deepseek-ai/dsh-session', async () => {
  146. const actual = await vi.importActual<typeof import('@deepseek-ai/dsh-session')>(
  147. '@deepseek-ai/dsh-session',
  148. )
  149. return {
  150. ...actual,
  151. SESSION_FORMAT_VERSION: currentVersion,
  152. Session: { create: validateHeader },
  153. }
  154. })
  155. vi.doMock('../src/format-migrations/index.ts', () => ({
  156. SESSION_FORMAT_MIGRATIONS: migrations,
  157. }))
  158. const decoder = await import('../src/format-decoder.ts')
  159. return {
  160. decodeStoredSession: decoder.decodeStoredSession,
  161. decodeStoredSessionHeader: decoder.decodeStoredSessionHeader,
  162. validateHeader,
  163. }
  164. }
  165. afterEach(() => {
  166. vi.doUnmock('@deepseek-ai/dsh-session')
  167. vi.doUnmock('../src/format-migrations/index.ts')
  168. vi.resetModules()
  169. })
  170. describe('versioned Session format decoder', { concurrent: false }, () => {
  171. it('describes both unsupported format directions', () => {
  172. expect(sessionFormatVersionRefusal(id, 1)).toContain('newer harness')
  173. expect(sessionFormatVersionRefusal(id, -1)).toContain('older than the supported')
  174. })
  175. it('runs a single migration lazily and reads the complete old log before slicing', async () => {
  176. const calls: string[] = []
  177. const step = migration(0, calls)
  178. const { decodeStoredSession, validateHeader } = await configuredDecoder(1, [step], calls)
  179. const originalEvents = eventLog()
  180. const originalSnapshot = structuredClone(originalEvents)
  181. const stored = storedSource(0, originalEvents)
  182. const decoded = decodeStoredSession(stored.source, id, 1)
  183. expect(decoded.sourceVersion).toBe(0)
  184. expect(decoded.meta.version).toBe(1)
  185. expect(calls).toEqual(['header:0', 'validate-header'])
  186. expect(stored.reads).toEqual([])
  187. const migrated = await collectEvents(decoded.events)
  188. await decoded.completed
  189. expect(stored.reads).toEqual([0])
  190. expect(calls).toEqual(['header:0', 'validate-header', 'events:0'])
  191. expect(migrated).toEqual([
  192. {
  193. ...originalEvents[1],
  194. data: { ...originalEvents[1]?.data, migratedFrom0: true, migrationPath: [0] },
  195. },
  196. ])
  197. expect(originalEvents).toEqual(originalSnapshot)
  198. expect(stored.meta).toEqual({ version: 0, id, createdAt: 1 })
  199. expect(validateHeader).toHaveBeenCalledOnce()
  200. })
  201. it('lets an old-format suffix migration use facts from events before fromSeq', async () => {
  202. const step = defineMigration(0, () => {
  203. let previousSeq: number | undefined
  204. return {
  205. header: meta => ({ ...(meta as Record<string, unknown>), version: 1 }),
  206. event(value) {
  207. const event = value as SessionEvent
  208. const migrated = previousSeq === undefined
  209. ? event
  210. : { ...event, data: { ...event.data, previousSeq } }
  211. previousSeq = event.seq
  212. return migrated
  213. },
  214. }
  215. })
  216. const { decodeStoredSession } = await configuredDecoder(1, [step])
  217. const stored = storedSource(0, eventLog())
  218. const decoded = decodeStoredSession(stored.source, id, 1)
  219. const events = await collectEvents(decoded.events)
  220. await decoded.completed
  221. expect(stored.reads).toEqual([0])
  222. expect(events).toEqual([{
  223. ...eventLog()[1],
  224. data: { ...eventLog()[1]?.data, previousSeq: 0 },
  225. }])
  226. })
  227. it('streams migrated events with backpressure instead of buffering the complete log', async () => {
  228. const releaseTail = Promise.withResolvers<undefined>()
  229. const physicalCompletion = Promise.withResolvers<StoredEventReadCompletion<never>>()
  230. const reads: number[] = []
  231. const source: StoredSessionSource<never> = {
  232. meta: { version: 0, id, createdAt: 1 },
  233. revision: SessionPersistenceRevision('streaming-source'),
  234. readEvents({ fromSeq = 0 } = {}) {
  235. reads.push(fromSeq)
  236. return {
  237. events: (async function* (): AsyncIterable<unknown> {
  238. try {
  239. yield structuredClone(eventLog()[0])
  240. await releaseTail.promise
  241. yield structuredClone(eventLog()[1])
  242. physicalCompletion.resolve({})
  243. } catch (error: unknown) {
  244. physicalCompletion.reject(error)
  245. throw error
  246. }
  247. })(),
  248. completed: physicalCompletion.promise,
  249. }
  250. },
  251. }
  252. const { decodeStoredSession } = await configuredDecoder(1, [migration(0, [])])
  253. const decoded = decodeStoredSession(source, id)
  254. const iterator = decoded.events[Symbol.asyncIterator]()
  255. const first = await iterator.next()
  256. expect(first).toMatchObject({ done: false, value: { seq: 0 } })
  257. expect(reads).toEqual([0])
  258. let completed = false
  259. void decoded.completed.then(() => { completed = true })
  260. await Promise.resolve()
  261. expect(completed).toBe(false)
  262. releaseTail.resolve(undefined)
  263. await expect(iterator.next()).resolves.toMatchObject({ done: false, value: { seq: 1 } })
  264. await expect(iterator.next()).resolves.toEqual({ done: true, value: undefined })
  265. await expect(decoded.completed).resolves.toEqual({})
  266. })
  267. it('runs a complete multi-step chain before current header and event validation', async () => {
  268. const calls: string[] = []
  269. const { decodeStoredSession } = await configuredDecoder(
  270. 2,
  271. [migration(0, calls), migration(1, calls)],
  272. calls,
  273. )
  274. const stored = storedSource(0, eventLog())
  275. const decoded = decodeStoredSession(stored.source, id)
  276. expect(decoded.meta.version).toBe(2)
  277. expect(calls).toEqual(['header:0', 'header:1', 'validate-header'])
  278. const events = await collectEvents(decoded.events)
  279. await decoded.completed
  280. expect(calls).toEqual([
  281. 'header:0',
  282. 'header:1',
  283. 'validate-header',
  284. 'events:0',
  285. 'events:1',
  286. ])
  287. expect(events[0]?.data).toMatchObject({ migratedFrom0: true, migratedFrom1: true })
  288. expect(events[0]?.data).toMatchObject({ migrationPath: [0, 1] })
  289. })
  290. it('detaches each migration output before the next migration mutates its input', async () => {
  291. const retained: Array<Record<string, unknown>> = []
  292. const first = defineMigration(0, () => ({
  293. header: meta => ({ ...(meta as Record<string, unknown>), version: 1 }),
  294. event(value) {
  295. const event = value as SessionEvent
  296. const output = {
  297. ...event,
  298. data: { ...(event.data as Record<string, unknown>), first: true },
  299. }
  300. retained.push(output.data)
  301. return output
  302. },
  303. }))
  304. const second = defineMigration(1, () => ({
  305. header: meta => ({ ...(meta as Record<string, unknown>), version: 2 }),
  306. event(value) {
  307. const event = value as SessionEvent
  308. const data = event.data as Record<string, unknown>
  309. data['second'] = true
  310. return event
  311. },
  312. }))
  313. const { decodeStoredSession } = await configuredDecoder(2, [first, second])
  314. const decoded = decodeStoredSession(storedSource(0, eventLog()).source, id)
  315. const events = await collectEvents(decoded.events)
  316. await decoded.completed
  317. expect(events.every(event => (event.data as Record<string, unknown>)['second'] === true)).toBe(true)
  318. expect(retained.every(data => data['second'] === undefined)).toBe(true)
  319. })
  320. it('rejects a non-JSON event output before a later migration can repair it', async () => {
  321. const first = defineMigration(0, () => ({
  322. header: meta => ({ ...(meta as Record<string, unknown>), version: 1 }),
  323. event(value) {
  324. const event = value as SessionEvent
  325. return {
  326. ...event,
  327. data: { ...(event.data as Record<string, unknown>), transient: undefined },
  328. }
  329. },
  330. }))
  331. const second = defineMigration(1, () => ({
  332. header: meta => ({ ...(meta as Record<string, unknown>), version: 2 }),
  333. event(value) {
  334. const event = value as SessionEvent
  335. const data = event.data as Record<string, unknown>
  336. delete data['transient']
  337. return event
  338. },
  339. }))
  340. const { decodeStoredSession } = await configuredDecoder(2, [first, second])
  341. const failure = await decodedFailure(
  342. decodeStoredSession(storedSource(0, eventLog()).source, id),
  343. )
  344. expect(failure.message).toMatch(/event migration v0 -> v1 failed at seq 0/)
  345. expect((failure.cause as Error).message).toMatch(/not losslessly JSON-serializable/)
  346. })
  347. it('plans by version even when registry entries are declared out of order', async () => {
  348. const calls: string[] = []
  349. const { decodeStoredSession } = await configuredDecoder(
  350. 2,
  351. [migration(1, calls), migration(0, calls)],
  352. calls,
  353. )
  354. const decoded = decodeStoredSession(storedSource(0, eventLog()).source, id)
  355. const events = await collectEvents(decoded.events)
  356. await decoded.completed
  357. expect(calls.slice(0, 3)).toEqual(['header:0', 'header:1', 'validate-header'])
  358. expect(events[0]?.data).toMatchObject({ migrationPath: [0, 1] })
  359. })
  360. it('starts a multi-version registry at the source version', async () => {
  361. const calls: string[] = []
  362. const { decodeStoredSession } = await configuredDecoder(
  363. 2,
  364. [migration(0, calls), migration(1, calls)],
  365. calls,
  366. )
  367. const stored = storedSource(1, eventLog())
  368. const decoded = decodeStoredSession(stored.source, id, 1)
  369. const events = await collectEvents(decoded.events)
  370. await decoded.completed
  371. expect(calls).toEqual(['header:1', 'validate-header', 'events:1'])
  372. expect(stored.reads).toEqual([0])
  373. expect(events[0]?.data).toMatchObject({ migrationPath: [1] })
  374. })
  375. it('retains instance state from the header through events and finishes at EOF', async () => {
  376. const calls: string[] = []
  377. const Migration = defineMigration(0, () => {
  378. let headerId: SessionId | undefined
  379. let migratedEvents = 0
  380. return {
  381. header(meta) {
  382. calls.push('header')
  383. headerId = SessionId((meta as Record<string, unknown>)['id'] as string)
  384. return { ...(meta as Record<string, unknown>), version: 1 }
  385. },
  386. event(value) {
  387. calls.push(`event:${migratedEvents}`)
  388. migratedEvents += 1
  389. return {
  390. ...(value as SessionEvent),
  391. data: { ...(value as SessionEvent).data, headerId, migratedEvents },
  392. }
  393. },
  394. finish() {
  395. calls.push(`finish:${migratedEvents}`)
  396. },
  397. }
  398. })
  399. const { decodeStoredSession } = await configuredDecoder(1, [Migration])
  400. const decoded = decodeStoredSession(storedSource(0, eventLog()).source, id)
  401. expect(calls).toEqual(['header'])
  402. const events = await collectEvents(decoded.events)
  403. await decoded.completed
  404. expect(calls).toEqual(['header', 'event:0', 'event:1', 'finish:2'])
  405. expect(events.map(event => event.data)).toMatchObject([
  406. { headerId: id, migratedEvents: 1 },
  407. { headerId: id, migratedEvents: 2 },
  408. ])
  409. })
  410. it('migrates and validates a header without requiring an event source', async () => {
  411. const calls: string[] = []
  412. const first = defineMigration(0, () => ({
  413. header(meta) {
  414. calls.push('header:0')
  415. return { ...(meta as Record<string, unknown>), version: 1 }
  416. },
  417. event: value => value,
  418. finish() {
  419. calls.push('finish:0')
  420. },
  421. }))
  422. const { decodeStoredSessionHeader } = await configuredDecoder(
  423. 2,
  424. [first, migration(1, calls)],
  425. calls,
  426. )
  427. const header = decodeStoredSessionHeader({ version: 0, id, createdAt: 1 }, id)
  428. expect(header.version).toBe(2)
  429. expect(calls).toEqual(['header:0', 'header:1', 'validate-header'])
  430. })
  431. it('allows a migration instance without finish', async () => {
  432. class MigrationWithoutFinish implements SessionFormatMigrationInstance {
  433. static readonly from = 0
  434. static readonly to = 1
  435. header(meta: unknown): unknown {
  436. return { ...(meta as Record<string, unknown>), version: 1 }
  437. }
  438. event(value: unknown): unknown {
  439. return value
  440. }
  441. }
  442. const { decodeStoredSession } = await configuredDecoder(1, [MigrationWithoutFinish])
  443. const decoded = decodeStoredSession(storedSource(0, eventLog()).source, id)
  444. await expect(collectEvents(decoded.events)).resolves.toEqual(eventLog())
  445. await expect(decoded.completed).resolves.toEqual({})
  446. })
  447. it('applies event migration before the current event vocabulary check', async () => {
  448. const step = defineMigration(0, () => ({
  449. header: meta => ({ ...(meta as Record<string, unknown>), version: 1 }),
  450. event(value) {
  451. const event = value as Record<string, unknown>
  452. return { ...event, type: 'turn/start', data: { turn: 1 } }
  453. },
  454. }))
  455. const { decodeStoredSession } = await configuredDecoder(1, [step])
  456. const stored = storedSource(0, [
  457. { type: 'legacy/turn-begin', seq: 0, time: 1, data: { legacyTurn: 1 } },
  458. ])
  459. const decoded = decodeStoredSession(stored.source, id)
  460. const events = await collectEvents(decoded.events)
  461. await decoded.completed
  462. expect(events).toEqual([
  463. { type: 'turn/start', seq: 0, time: 1, data: { turn: 1 } },
  464. ])
  465. })
  466. it('uses suffix access directly for the current format', async () => {
  467. const { decodeStoredSession } = await configuredDecoder(2, [])
  468. const stored = storedSource(2, eventLog())
  469. const decoded = decodeStoredSession(stored.source, id, 1)
  470. const events = await collectEvents(decoded.events)
  471. await decoded.completed
  472. expect(stored.reads).toEqual([1])
  473. expect(events).toEqual(eventLog().slice(1))
  474. })
  475. it('buffers a safe current-v0 suffix once without reopening the prefix', async () => {
  476. const { decodeStoredSession } = await configuredDecoder(0, [])
  477. const stored = storedSource(0, eventLog())
  478. const decoded = decodeStoredSession(stored.source, id, 1)
  479. expect(await collectEvents(decoded.events)).toEqual(eventLog().slice(1))
  480. await expect(decoded.completed).resolves.toEqual({})
  481. expect(stored.reads).toEqual([1])
  482. })
  483. it('reopens the complete current-v0 log when a legacy suffix record needs its prefix', async () => {
  484. const { decodeStoredSession } = await configuredDecoder(0, [])
  485. const legacy = [
  486. { type: 'turn/start', seq: 0, time: 1, data: { turn: 1 } },
  487. {
  488. type: 'steering/message',
  489. seq: 1,
  490. time: 2,
  491. data: { turn: 1, content: [{ type: 'text', text: 'continue' }], source: { kind: 'user' } },
  492. },
  493. ]
  494. const stored = storedSource(0, legacy)
  495. const decoded = decodeStoredSession(stored.source, id, 1)
  496. const events = await collectEvents(decoded.events)
  497. await decoded.completed
  498. expect(stored.reads).toEqual([1, 0])
  499. expect(events).toMatchObject([{ type: 'user/message', seq: 1 }])
  500. })
  501. it('observes a failed physical completion after reopening a required v0 prefix', async () => {
  502. const failure = new SessionPersistenceRevisionConflictError('reopened prefix changed')
  503. const fullCompletion = Promise.withResolvers<StoredEventReadCompletion<never>>()
  504. const source: StoredSessionSource<never> = {
  505. meta: { version: 0, id, createdAt: 1 },
  506. revision: SessionPersistenceRevision('prefix-conflict'),
  507. readEvents({ fromSeq = 0 } = {}) {
  508. if (fromSeq > 0) {
  509. return {
  510. events: (async function* (): AsyncIterable<unknown> {
  511. yield {
  512. type: 'steering/message', seq: 1, time: 2,
  513. data: { turn: 1, content: [], source: { kind: 'user' } },
  514. }
  515. })(),
  516. completed: Promise.resolve({}),
  517. }
  518. }
  519. return {
  520. events: (async function* (): AsyncIterable<unknown> {
  521. fullCompletion.reject(failure)
  522. throw failure
  523. })(),
  524. completed: fullCompletion.promise,
  525. }
  526. },
  527. }
  528. const { decodeStoredSession } = await configuredDecoder(0, [])
  529. const decoded = decodeStoredSession(source, id, 1)
  530. await expect(decodedFailure(decoded)).resolves.toBe(failure)
  531. })
  532. it('classifies every v0 prefix-independent suffix value without assuming a record', () => {
  533. const compatibility = unversionedFormatCompatibility(0)
  534. if (compatibility === undefined) throw new Error('v0 compatibility must be registered')
  535. expect(compatibility.requiresPrefix(null)).toBe(false)
  536. expect(compatibility.requiresPrefix({ type: 'turn/end', data: null })).toBe(false)
  537. expect(compatibility.requiresPrefix({ type: 'user/message', data: { id: 'current', content: [] } })).toBe(false)
  538. expect(compatibility.requiresPrefix({ type: 'user/message', data: { content: [] } })).toBe(true)
  539. expect(compatibility.requiresPrefix({ type: 'assistant/message', data: { content: [] } })).toBe(true)
  540. expect(compatibility.requiresPrefix({ type: 'tool/result', data: { callId: 'call' } })).toBe(true)
  541. })
  542. it('preserves already-canonical v0 turn-end reasons', async () => {
  543. const compatibility = unversionedFormatCompatibility(0)
  544. if (compatibility === undefined) throw new Error('v0 compatibility must be registered')
  545. const events = [
  546. {
  547. type: 'turn/end', seq: 0, time: 1,
  548. data: { turn: 1, reason: { kind: 'aborted', reason: { kind: 'disposed' } } },
  549. },
  550. {
  551. type: 'turn/end', seq: 1, time: 2,
  552. data: { turn: 2, reason: { kind: 'error', error: { message: 'failed', code: 'UNKNOWN' } } },
  553. },
  554. ]
  555. const input = (async function* (): AsyncIterable<unknown> {
  556. yield* events
  557. })()
  558. const canonical: unknown[] = []
  559. for await (const event of compatibility.canonicalizeEvents(input, id)) canonical.push(event)
  560. expect(canonical).toEqual(events)
  561. })
  562. it('canonicalizes every historical compact event name without changing its record', async () => {
  563. const { decodeStoredSession } = await configuredDecoder(0, [])
  564. const events = [
  565. {
  566. type: 'compact/start', seq: 0, time: 1,
  567. data: { compactionId: 'legacy', turn: 1 },
  568. surfaceOp: { op: 'retain' },
  569. },
  570. {
  571. type: 'compact/summary', seq: 1, time: 2,
  572. data: { summary: 'old summary', shadowedSeqs: [7, 8] },
  573. durableMetadata: { source: 'historical-v0' },
  574. },
  575. {
  576. type: 'compaction/end', seq: 2, time: 3,
  577. data: { compactionId: 'current', turn: 1 },
  578. },
  579. {
  580. type: 'compact/end', seq: 3, time: 4,
  581. data: { compactionId: 'legacy', turn: 1 },
  582. },
  583. {
  584. type: 'compact/prune', seq: 4, time: 5,
  585. data: {
  586. shadowedRange: { start: 7, end: 8 },
  587. shadowedSeqs: [7, 8],
  588. shadowedTokenCount: 456,
  589. },
  590. },
  591. ]
  592. const stored = storedSource(0, events)
  593. const decoded = decodeStoredSession(stored.source, id)
  594. const canonical = await collectEvents(decoded.events)
  595. await decoded.completed
  596. expect(canonical).toEqual(events.map(event => ({
  597. ...event,
  598. type: event.type.replace(/^compact\//, 'compaction/'),
  599. })))
  600. expect(stored.reads).toEqual([0])
  601. })
  602. it('still rejects other unknown v0 event names after compaction normalization', async () => {
  603. const { decodeStoredSession } = await configuredDecoder(0, [])
  604. const decoded = decodeStoredSession(storedSource(0, [{
  605. type: 'compact/future', seq: 0, time: 1, data: {},
  606. }]).source, id)
  607. const failure = await decodedFailure(decoded)
  608. expect(failure.message).toMatch(/event type "compact\/future".*not marked ignorable/)
  609. })
  610. it('does not run older registered steps for an already-current source', async () => {
  611. const calls: string[] = []
  612. const { decodeStoredSession } = await configuredDecoder(
  613. 2,
  614. [migration(0, calls), migration(1, calls)],
  615. calls,
  616. )
  617. const stored = storedSource(2, eventLog())
  618. const decoded = decodeStoredSession(stored.source, id, 1)
  619. await collectEvents(decoded.events)
  620. await decoded.completed
  621. expect(calls).toEqual(['validate-header'])
  622. expect(stored.reads).toEqual([1])
  623. })
  624. it('opens a fresh revision-bound reader for each decode of the same source', async () => {
  625. const { decodeStoredSession } = await configuredDecoder(0, [])
  626. const stored = storedSource(0, eventLog())
  627. const first = decodeStoredSession(stored.source, id)
  628. expect(await collectEvents(first.events)).toEqual(eventLog())
  629. await first.completed
  630. const second = decodeStoredSession(stored.source, id)
  631. expect(await collectEvents(second.events)).toEqual(eventLog())
  632. await second.completed
  633. expect(first.revision).toBe(second.revision)
  634. expect(stored.reads).toEqual([0, 0])
  635. })
  636. it('propagates a physical revision conflict unchanged through events and completion', async () => {
  637. const failure = new SessionPersistenceRevisionConflictError('source changed')
  638. const physicalCompletion = Promise.withResolvers<StoredEventReadCompletion<never>>()
  639. const source: StoredSessionSource<never> = {
  640. meta: { version: 0, id, createdAt: 1 },
  641. revision: SessionPersistenceRevision('conflicting-source'),
  642. readEvents: () => ({
  643. events: (async function* (): AsyncIterable<unknown> {
  644. physicalCompletion.reject(failure)
  645. throw failure
  646. })(),
  647. completed: physicalCompletion.promise,
  648. }),
  649. }
  650. const { decodeStoredSession } = await configuredDecoder(0, [])
  651. const decoded = decodeStoredSession(source, id)
  652. const completion = decoded.completed.catch((error: unknown) => error)
  653. const consumption = collectEvents(decoded.events).catch((error: unknown) => error)
  654. const [streamFailure, completionFailure] = await Promise.all([consumption, completion])
  655. expect(streamFailure).toBe(failure)
  656. expect(completionFailure).toBe(failure)
  657. })
  658. it('propagates an upstream revision conflict unchanged through a migration step', async () => {
  659. const { decodeStoredSession } = await configuredDecoder(1, [migration(0, [])])
  660. const { SessionPersistenceRevisionConflictError: DecoderRevisionConflictError } = await import('../src/revision.ts')
  661. const failure = new DecoderRevisionConflictError('migrating source changed')
  662. const source: StoredSessionSource<never> = {
  663. meta: { version: 0, id, createdAt: 1 },
  664. revision: SessionPersistenceRevision('conflicting-migration-source'),
  665. readEvents: () => ({
  666. events: (async function* (): AsyncIterable<unknown> {
  667. throw failure
  668. })(),
  669. completed: Promise.reject(failure),
  670. }),
  671. }
  672. await expect(decodedFailure(decodeStoredSession(source, id))).resolves.toBe(failure)
  673. })
  674. it('rejects a missing path and a future source in the correct direction', async () => {
  675. const { decodeStoredSession, validateHeader } = await configuredDecoder(2, [])
  676. const old = storedSource(0, [])
  677. const future = storedSource(3, [])
  678. expect(() => decodeStoredSession(old.source, id))
  679. .toThrow(/missing v0 -> v1/)
  680. expect(() => decodeStoredSession(future.source, id))
  681. .toThrow(/newer harness/)
  682. expect(old.reads).toEqual([])
  683. expect(future.reads).toEqual([])
  684. expect(validateHeader).not.toHaveBeenCalled()
  685. })
  686. it('preserves the raw location in unsupported-format diagnostics', async () => {
  687. const { decodeStoredSession } = await configuredDecoder(0, [])
  688. const stored = storedSource(1, [])
  689. const location = { kind: 'jsonl', path: '/tmp/session.jsonl' }
  690. const source: StoredSessionSource<never> = { ...stored.source, location }
  691. let failure: unknown
  692. try {
  693. decodeStoredSession(source, id)
  694. } catch (error: unknown) {
  695. failure = error
  696. }
  697. expect(failure).toMatchObject({
  698. name: 'SessionFormatUnsupportedError',
  699. location,
  700. })
  701. expect((failure as Error).message).toContain('(raw log: /tmp/session.jsonl)')
  702. })
  703. it('rejects an invalid suffix before validating the header or opening events', async () => {
  704. const { decodeStoredSession, validateHeader } = await configuredDecoder(0, [])
  705. const stored = storedSource(0, eventLog())
  706. for (const fromSeq of [-1, 1.5, Number.MAX_SAFE_INTEGER + 1]) {
  707. expect(() => decodeStoredSession(stored.source, id, fromSeq))
  708. .toThrow(/fromSeq must be a non-negative safe integer/)
  709. }
  710. expect(validateHeader).not.toHaveBeenCalled()
  711. expect(stored.reads).toEqual([])
  712. })
  713. it('validates unknown durable header fields before path selection', async () => {
  714. const { decodeStoredSession, validateHeader } = await configuredDecoder(0, [])
  715. const cases: Array<{ meta: unknown; message: RegExp }> = [
  716. { meta: null, message: /header is not a lossless JSON record/ },
  717. { meta: { version: '0', id }, message: /invalid format version/ },
  718. { meta: { version: 0, id: 42 }, message: /has no string id/ },
  719. ]
  720. let reads = 0
  721. for (const entry of cases) {
  722. const source: StoredSessionSource<never> = {
  723. meta: entry.meta,
  724. revision: SessionPersistenceRevision('invalid-header'),
  725. readEvents: () => {
  726. reads += 1
  727. return { events: (async function* () {})(), completed: Promise.resolve({}) }
  728. },
  729. }
  730. expect(() => decodeStoredSession(source, id)).toThrow(entry.message)
  731. }
  732. expect(reads).toBe(0)
  733. expect(validateHeader).not.toHaveBeenCalled()
  734. })
  735. it('rejects every malformed current event envelope through the stream and completion', async () => {
  736. const { decodeStoredSession } = await configuredDecoder(1, [])
  737. const cases: Array<{ value: unknown; message: RegExp }> = [
  738. { value: null, message: /non-record event/ },
  739. { value: { seq: 0, time: 1, data: {} }, message: /without a string type/ },
  740. { value: { type: 'turn/start', seq: -1, time: 1, data: {} }, message: /invalid seq -1/ },
  741. { value: { type: 'turn/start', seq: 0, time: 'now', data: {} }, message: /invalid time/ },
  742. { value: { type: 'turn/start', seq: 0, time: 1 }, message: /without data/ },
  743. ]
  744. for (const entry of cases) {
  745. const decoded = decodeStoredSession(storedSource(1, [entry.value]).source, id)
  746. expect((await decodedFailure(decoded)).message).toMatch(entry.message)
  747. }
  748. })
  749. it('rejects a stored event that cannot be represented as JSON', async () => {
  750. const { decodeStoredSession } = await configuredDecoder(1, [])
  751. const decoded = decodeStoredSession(storedSource(1, [undefined]).source, id)
  752. expect((await decodedFailure(decoded)).message).toMatch(/not losslessly JSON-serializable/)
  753. })
  754. it('rejects every malformed v0 event before same-version canonicalization', async () => {
  755. const { decodeStoredSession } = await configuredDecoder(0, [])
  756. const cases: Array<{ value: unknown; message: RegExp }> = [
  757. { value: null, message: /non-record event/ },
  758. { value: { seq: 0, time: 1, data: {} }, message: /without a string type/ },
  759. { value: { type: 'turn/start', seq: -1, time: 1, data: {} }, message: /invalid seq -1/ },
  760. { value: { type: 'turn/start', seq: 0, time: 'now', data: {} }, message: /invalid time/ },
  761. { value: { type: 'turn/start', seq: 0, time: 1 }, message: /without data/ },
  762. ]
  763. for (const entry of cases) {
  764. const decoded = decodeStoredSession(storedSource(0, [entry.value]).source, id)
  765. expect((await decodedFailure(decoded)).message).toMatch(entry.message)
  766. }
  767. })
  768. it('lets a v0 turn/end with opaque data reach current validation unchanged', async () => {
  769. const { decodeStoredSession } = await configuredDecoder(0, [])
  770. const decoded = decodeStoredSession(storedSource(0, [
  771. { type: 'turn/end', seq: 0, time: 1, data: null },
  772. ]).source, id)
  773. await expect(collectEvents(decoded.events)).resolves.toEqual([
  774. { type: 'turn/end', seq: 0, time: 1, data: null },
  775. ])
  776. await expect(decoded.completed).resolves.toEqual({})
  777. })
  778. it('rejects a migration that returns the wrong header version', async () => {
  779. const bad = defineMigration(0, () => ({
  780. header: meta => ({ ...(meta as Record<string, unknown>), version: 0 }),
  781. event: value => value,
  782. }))
  783. const first = await configuredDecoder(1, [bad])
  784. expect(() => first.decodeStoredSession(storedSource(0, []).source, id))
  785. .toThrow(/returned header version 0/)
  786. expect(first.validateHeader).not.toHaveBeenCalled()
  787. const calls: string[] = []
  788. const badSecond = defineMigration(1, () => ({
  789. header: meta => ({ ...(meta as Record<string, unknown>), version: 1 }),
  790. event: value => value,
  791. }))
  792. const second = await configuredDecoder(2, [migration(0, calls), badSecond], calls)
  793. const stored = storedSource(0, [])
  794. expect(() => second.decodeStoredSession(stored.source, id))
  795. .toThrow(/v1 -> v2 returned header version 1/)
  796. expect(calls).toEqual(['header:0'])
  797. expect(second.validateHeader).not.toHaveBeenCalled()
  798. expect(stored.reads).toEqual([])
  799. })
  800. it('rejects a migration that changes the session id or cwd storage identity', async () => {
  801. const changedId = defineMigration(0, () => ({
  802. header: meta => ({ ...(meta as Record<string, unknown>), version: 1, id: 'other' }),
  803. event: value => value,
  804. }))
  805. const first = await configuredDecoder(1, [changedId])
  806. expect(() => first.decodeStoredSession(storedSource(0, []).source, id))
  807. .toThrow(/changed session storage identity/)
  808. const changedCwd = defineMigration(0, () => ({
  809. header: meta => ({ ...(meta as Record<string, unknown>), version: 1, cwd: '/other' }),
  810. event: value => value,
  811. }))
  812. const second = await configuredDecoder(1, [changedCwd])
  813. const stored = storedSource(0, [])
  814. stored.meta['cwd'] = '/work'
  815. expect(() => second.decodeStoredSession(stored.source, id))
  816. .toThrow(/changed session storage identity/)
  817. })
  818. it('wraps a header migration failure with the failing version step', async () => {
  819. const cause = new Error('bad legacy header')
  820. const step = defineMigration(0, () => ({
  821. header: () => { throw cause },
  822. event: value => value,
  823. }))
  824. const { decodeStoredSession } = await configuredDecoder(1, [step])
  825. let failure: unknown
  826. try {
  827. decodeStoredSession(storedSource(0, []).source, id)
  828. } catch (error: unknown) {
  829. failure = error
  830. }
  831. expect(failure).toMatchObject({
  832. message: `session "${id}" header migration v0 -> v1 failed`,
  833. cause,
  834. })
  835. })
  836. it('mirrors an event migration failure through the stream and completion promise', async () => {
  837. const cause = new Error('bad legacy event')
  838. const step = defineMigration(0, () => ({
  839. header: meta => ({ ...(meta as Record<string, unknown>), version: 1 }),
  840. event: () => { throw cause },
  841. }))
  842. const { decodeStoredSession } = await configuredDecoder(1, [step])
  843. const decoded = decodeStoredSession(storedSource(0, eventLog()).source, id)
  844. const completion = decoded.completed.catch((error: unknown) => error)
  845. const consumption = collectEvents(decoded.events).catch((error: unknown) => error)
  846. const [streamFailure, completionFailure] = await Promise.all([consumption, completion])
  847. expect(streamFailure).toBe(completionFailure)
  848. expect(streamFailure).toMatchObject({
  849. message: `session "${id}" event migration v0 -> v1 failed at seq 0`,
  850. cause,
  851. })
  852. })
  853. it('mirrors a finish failure through the stream and completion promise', async () => {
  854. const cause = new Error('unclosed legacy state')
  855. const Migration = defineMigration(0, () => ({
  856. header: meta => ({ ...(meta as Record<string, unknown>), version: 1 }),
  857. event: value => value,
  858. finish: () => { throw cause },
  859. }))
  860. const { decodeStoredSession } = await configuredDecoder(1, [Migration])
  861. const failure = await decodedFailure(decodeStoredSession(storedSource(0, eventLog()).source, id))
  862. expect(failure).toMatchObject({
  863. message: `session "${id}" event migration v0 -> v1 failed at EOF`,
  864. cause,
  865. })
  866. })
  867. it('rejects a migration that changes an event sequence number', async () => {
  868. const step = defineMigration(0, () => ({
  869. header: meta => ({ ...(meta as Record<string, unknown>), version: 1 }),
  870. event(value) {
  871. const event = value as SessionEvent
  872. return { ...event, seq: event.seq + 1 }
  873. },
  874. }))
  875. const { decodeStoredSession } = await configuredDecoder(1, [step])
  876. const decoded = decodeStoredSession(storedSource(0, eventLog()).source, id)
  877. const completion = decoded.completed.catch((error: unknown) => error)
  878. const consumption = collectEvents(decoded.events).catch((error: unknown) => error)
  879. const [streamFailure, completionFailure] = await Promise.all([consumption, completion])
  880. expect(streamFailure).toBe(completionFailure)
  881. expect((streamFailure as Error).message).toMatch(/changed event seq 0 to 1/)
  882. })
  883. it('rejects a non-contiguous current-format event sequence', async () => {
  884. const { decodeStoredSession } = await configuredDecoder(0, [])
  885. const stored = storedSource(0, [
  886. { type: 'turn/start', seq: 1, time: 1, data: { turn: 1 } },
  887. ])
  888. const failure = await decodedFailure(decodeStoredSession(stored.source, id))
  889. expect(failure.message).toContain(`session "${id}" event seq mismatch: expected 0, got 1`)
  890. })
  891. it('runs current header validation only after the final header step', async () => {
  892. const calls: string[] = []
  893. const finalStep = defineMigration(1, () => ({
  894. header(meta) {
  895. calls.push('header:1')
  896. const { createdAt: _createdAt, ...rest } = meta as Record<string, unknown>
  897. return { ...rest, version: 2 }
  898. },
  899. event: value => value,
  900. }))
  901. const { decodeStoredSession } = await configuredDecoder(
  902. 2,
  903. [migration(0, calls), finalStep],
  904. calls,
  905. )
  906. const stored = storedSource(0, eventLog())
  907. expect(() => decodeStoredSession(stored.source, id))
  908. .toThrow(/current header validator received invalid createdAt/)
  909. expect(calls).toEqual(['header:0', 'header:1', 'validate-header'])
  910. expect(stored.reads).toEqual([])
  911. })
  912. it('detaches stored header and event objects before a mutating migration runs', async () => {
  913. const originalEvents = eventLog()
  914. const eventSnapshot = structuredClone(originalEvents)
  915. const step = defineMigration(0, () => ({
  916. header(meta) {
  917. const record = meta as Record<string, unknown>
  918. record['version'] = 1
  919. return record
  920. },
  921. event(value) {
  922. const event = value as SessionEvent
  923. const data = event.data as Record<string, unknown>
  924. data['mutated'] = true
  925. return event
  926. },
  927. }))
  928. const { decodeStoredSession } = await configuredDecoder(1, [step])
  929. const stored = storedSource(0, originalEvents)
  930. const decoded = decodeStoredSession(stored.source, id)
  931. const migrated = await collectEvents(decoded.events)
  932. await decoded.completed
  933. expect(migrated.every(event => (event.data as Record<string, unknown>)['mutated'] === true)).toBe(true)
  934. expect(stored.meta).toEqual({ version: 0, id, createdAt: 1 })
  935. expect(originalEvents).toEqual(eventSnapshot)
  936. })
  937. it('rejects duplicate, invalid, and future-targeting static registries at initialization', async () => {
  938. const calls: string[] = []
  939. await expect(configuredDecoder(1, [migration(0, calls), migration(0, calls)]))
  940. .rejects.toThrow(/duplicate Session format migration/)
  941. await expect(configuredDecoder(1, [migration(-1, calls)]))
  942. .rejects.toThrow(/adjacent non-negative version/)
  943. const nonAdjacent = migration(0, calls, 2)
  944. await expect(configuredDecoder(2, [nonAdjacent]))
  945. .rejects.toThrow(/adjacent non-negative version/)
  946. const fractional = migration(0.5, calls, 1.5)
  947. await expect(configuredDecoder(2, [fractional]))
  948. .rejects.toThrow(/adjacent non-negative version/)
  949. await expect(configuredDecoder(1, [migration(1, calls)]))
  950. .rejects.toThrow(/targets a version newer than this build/)
  951. })
  952. it('initializes with a gapped registry and refuses only sessions at or below the gap', async () => {
  953. const calls: string[] = []
  954. const { decodeStoredSession } = await configuredDecoder(
  955. 3,
  956. [migration(0, calls), migration(2, calls)],
  957. calls,
  958. )
  959. const current = storedSource(3, [])
  960. const pastGap = storedSource(2, eventLog())
  961. const atGap = storedSource(1, eventLog())
  962. const belowGap = storedSource(0, eventLog())
  963. expect(decodeStoredSession(current.source, id).meta.version).toBe(3)
  964. const decoded = decodeStoredSession(pastGap.source, id)
  965. expect(decoded.sourceVersion).toBe(2)
  966. expect(decoded.meta.version).toBe(3)
  967. expect(() => decodeStoredSession(atGap.source, id))
  968. .toThrow(/missing v1 -> v2/)
  969. expect(() => decodeStoredSession(belowGap.source, id))
  970. .toThrow(/missing v1 -> v2/)
  971. expect(pastGap.reads).toEqual([])
  972. })
  973. })