jsonl.spec.ts 63 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350
  1. import { MessageId, createUserMessage, createMessage } from '@deepseek-ai/dsh-llm'
  2. import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
  3. import { Context } from 'cordis'
  4. import { appendFile, mkdtemp, mkdir, rm, readFile, writeFile, readdir, stat, symlink } from 'node:fs/promises'
  5. import { tmpdir } from 'node:os'
  6. import { isAbsolute, join, relative, resolve } from 'node:path'
  7. import SessionStore, { SessionId } from '@deepseek-ai/dsh-session'
  8. import type { Session, SessionEvent, SessionHeader } from '@deepseek-ai/dsh-session'
  9. import SessionPersistenceJsonl from '@deepseek-ai/dsh-session-persistence-jsonl'
  10. import {
  11. encodeSegment, eventLines, logPath, projectDir, projectKey, scanLog, sessionDir, toHeaderLine,
  12. } from '../src/format.ts'
  13. import { runPersistenceContract, meta, oneTurnLog, appendLog } from '../../session-persistence/tests/contract.ts'
  14. import { runCoordinatorContract, type CoordinatorFixture } from '../../session-persistence/tests/coordinator-contract.ts'
  15. let root: string
  16. const dirs: string[] = []
  17. type MutableSessionHeader = { -readonly [K in keyof SessionHeader]: SessionHeader[K] }
  18. /** Test-only mutable view used to verify that backends detach returned/caller metadata. */
  19. function mutableHeader(header: SessionHeader): MutableSessionHeader {
  20. return header
  21. }
  22. /** Rewrite only a stored header while preserving every event byte below it. */
  23. async function rewriteHeader(path: string, update: (header: Record<string, unknown>) => void): Promise<void> {
  24. const lines = (await readFile(path, 'utf8')).split('\n')
  25. const header = JSON.parse(lines[0] as string) as Record<string, unknown>
  26. update(header)
  27. lines[0] = JSON.stringify(header)
  28. await writeFile(path, lines.join('\n'))
  29. }
  30. async function expectFlushError(promise: Promise<unknown>, message: RegExp): Promise<void> {
  31. try {
  32. await promise
  33. } catch (error) {
  34. expect(error).toBeInstanceOf(Error)
  35. expect((error as Error).message).toMatch(message)
  36. return
  37. }
  38. throw new Error('expected flush to reject')
  39. }
  40. async function freshRoot(): Promise<string> {
  41. const dir = await mkdtemp(join(tmpdir(), 'dsh-jsonl-'))
  42. dirs.push(dir)
  43. return dir
  44. }
  45. function rawLogPath(root: string, cwd: string | undefined, id: SessionId): string {
  46. return logPath(root, cwd, id, 'none')
  47. }
  48. afterEach(async () => {
  49. vi.restoreAllMocks()
  50. for (const d of dirs.splice(0)) await rm(d, { recursive: true, force: true })
  51. })
  52. function appendClosedTurn(session: Session): void {
  53. session.append('turn/start', { turn: 1 })
  54. session.append('user/message', createUserMessage({
  55. content: [{ type: 'text', text: 'hello' }],
  56. source: { kind: 'user' },
  57. }), { surfaceOp: 'append' })
  58. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  59. }
  60. // Run the shared backend contract against the real JSONL backend.
  61. runPersistenceContract('jsonl-none', async () => {
  62. const dir = await mkdtemp(join(tmpdir(), 'dsh-jsonl-'))
  63. const ctx = new Context()
  64. await ctx.plugin(SessionStore)
  65. const fiber = await ctx.plugin(SessionPersistenceJsonl, { root: dir, compression: 'none' })
  66. return {
  67. persistence: ctx.sessionPersistence,
  68. dispose: async () => {
  69. await fiber.dispose()
  70. await rm(dir, { recursive: true, force: true })
  71. },
  72. }
  73. })
  74. // Two mounts share this temp root to exercise reload. `corruptTail` appends a partial,
  75. // newline-less fragment past the committed region so coordinator repair runs on real file bytes.
  76. runCoordinatorContract('jsonl-none', async (): Promise<CoordinatorFixture> => {
  77. const dir = await mkdtemp(join(tmpdir(), 'dsh-jsonl-coord-'))
  78. return {
  79. mount: async (ctx) => {
  80. const fiber = await ctx.plugin(SessionPersistenceJsonl, { root: dir, compression: 'none' })
  81. return fiber
  82. },
  83. corruptTail: async (id, cwd) => {
  84. // A half-written record with no trailing newline: scanLog treats it as an
  85. // uncommitted crash fragment and reports committedBytes < byteLength, so
  86. // the coordinator sees a tornMarker to truncate.
  87. await appendFile(rawLogPath(dir, cwd, id), '{"type":"assistant/chunk","seq":8,"ti')
  88. },
  89. cleanup: async () => { await rm(dir, { recursive: true, force: true }) },
  90. }
  91. })
  92. describe('SessionPersistenceJsonl: format helpers', () => {
  93. it('encodeSegment neutralizes traversal, separators, and absolute paths', () => {
  94. expect(encodeSegment('..')).toBe('~002E~002E')
  95. expect(encodeSegment('.')).toBe('~002E')
  96. expect(encodeSegment('a/b')).toBe('a~002Fb')
  97. expect(encodeSegment('/etc/passwd')).toBe('~002Fetc~002Fpasswd')
  98. expect(encodeSegment('a\u0000b')).toBe('a~0000b')
  99. expect(encodeSegment('plain-ID_1.2')).toBe('plain-ID_1.2') // safe chars pass through
  100. expect(encodeSegment('a~b')).toBe('a~007Eb') // ~ itself is escaped
  101. })
  102. it('encodeSegment is injective over UTF-16, incl. lone surrogates', () => {
  103. // Distinct lone surrogates must NOT collide (Buffer.from would normalize
  104. // both to U+FFFD; code-unit escaping keeps them distinct).
  105. const hi = encodeSegment(String.fromCharCode(0xD800))
  106. const lo = encodeSegment(String.fromCharCode(0xDC00))
  107. expect(hi).toBe('~D800')
  108. expect(lo).toBe('~DC00')
  109. expect(hi).not.toBe(lo)
  110. // A literal "~002F" input cannot collide with the encoding of "/".
  111. expect(encodeSegment('~002F')).not.toBe(encodeSegment('/'))
  112. })
  113. it('encodeSegment rejects an empty id', () => {
  114. expect(() => encodeSegment('')).toThrow(/empty/)
  115. })
  116. it('projectKey normalizes project paths into bounded readable names', () => {
  117. expect(projectKey('/Users/qyj/work/deepseek-harness')).toBe('--Users-qyj-work-deepseek-harness--')
  118. expect(projectKey('/a/b-c')).toBe(projectKey('/a-b/c'))
  119. expect(projectKey('C:\\work\\agent')).toBe('--C-work-agent--')
  120. expect(projectKey('/开发/~agent')).toBe('--~5F00~53D1-~007Eagent--')
  121. expect(projectKey('/')).toBe('--root--')
  122. expect(projectKey('/' + 'x'.repeat(1_000))).toHaveLength(255)
  123. expect(() => projectKey('')).toThrow(/empty project path/)
  124. })
  125. it('resolves a relative custom root before locating a session', async () => {
  126. const absoluteRoot = await freshRoot()
  127. const ctx = new Context()
  128. await ctx.plugin(SessionStore)
  129. const fiber = await ctx.plugin(SessionPersistenceJsonl, {
  130. root: relative(process.cwd(), absoluteRoot),
  131. compression: 'none',
  132. })
  133. const m = meta('relative-location', '/work')
  134. expect(ctx.sessionPersistence.locate(m)).toEqual({
  135. kind: 'jsonl',
  136. path: rawLogPath(resolve(absoluteRoot), '/work', m.id),
  137. })
  138. await fiber.dispose()
  139. })
  140. })
  141. describe('SessionPersistenceJsonl: durability and crash semantics', () => {
  142. let ctx: Context
  143. beforeEach(async () => {
  144. root = await freshRoot()
  145. ctx = new Context()
  146. await ctx.plugin(SessionStore)
  147. await ctx.plugin(SessionPersistenceJsonl, { root, compression: 'none' })
  148. })
  149. afterEach(async () => { await ctx.fiber.dispose() })
  150. it('lazy materialization: create() writes no file until the first append', async () => {
  151. const m = meta('lazy', '/work')
  152. const location = ctx.sessionPersistence.locate(m)
  153. expect(location).toEqual({ kind: 'jsonl', path: rawLogPath(root, '/work', m.id) })
  154. expect(isAbsolute(location!.path)).toBe(true)
  155. await ctx.sessionPersistence.create(m)
  156. // locate() is a pure target-path calculation: neither it nor create()
  157. // materializes a file before the first append.
  158. const dir = sessionDir(root, '/work', m.id)
  159. await expect(stat(rawLogPath(root, '/work', m.id))).rejects.toThrow()
  160. expect((await ctx.sessionPersistence.list()).map(h => h.id)).not.toContain(m.id)
  161. await ctx.sessionPersistence.append(m.id, oneTurnLog())
  162. // now materialized
  163. expect((await stat(dir)).isDirectory()).toBe(true)
  164. expect((await stat(rawLogPath(root, '/work', m.id))).isFile()).toBe(true)
  165. expect((await ctx.sessionPersistence.list()).map(h => h.id)).toContain(m.id)
  166. })
  167. it('keeps the same location on resume and gives a fork its own location', async () => {
  168. const parent = meta('location-parent', '/work')
  169. const parentLocation = ctx.sessionPersistence.locate(parent)
  170. await ctx.sessionPersistence.create(parent)
  171. await ctx.sessionPersistence.append(parent.id, oneTurnLog())
  172. const loaded = await ctx.sessionPersistence.load(parent.id)
  173. expect(ctx.sessionPersistence.locate(loaded.meta)).toEqual(parentLocation)
  174. const child = {
  175. ...loaded.meta,
  176. id: SessionId('location-child'),
  177. parentSession: parent.id,
  178. seedLength: loaded.events.length,
  179. }
  180. const childLocation = ctx.sessionPersistence.locate(child)
  181. expect(childLocation?.path).not.toBe(parentLocation?.path)
  182. expect(childLocation).toEqual({ kind: 'jsonl', path: rawLogPath(root, '/work', child.id) })
  183. })
  184. it('round-trip is byte-identical (incl. assistant/chunk verbatim)', async () => {
  185. const m = meta('chunks')
  186. const log: SessionEvent[] = [
  187. { type: 'turn/start', seq: 0, time: 1, data: { turn: 1 } },
  188. { type: 'step/start', seq: 1, time: 2, data: { turn: 1, step: 1 } },
  189. { type: 'assistant/chunk', seq: 2, time: 3, data: { turn: 1, step: 1, chunk: { type: 'text-delta', index: 0, text: 'he' } } },
  190. { type: 'assistant/chunk', seq: 3, time: 4, data: { turn: 1, step: 1, chunk: { type: 'text-delta', index: 0, text: 'llo' } } },
  191. { type: 'assistant/message', seq: 4, time: 5, data: {
  192. turn: 1, step: 1,
  193. message: createMessage({
  194. role: 'assistant',
  195. content: [{ type: 'text', text: 'hello' }],
  196. source: {
  197. kind: 'model',
  198. ...{ provider: 'mock', model: 'mock' },
  199. },
  200. }),
  201. }, surfaceOp: 'append', sourceEventSeqs: [2, 3] },
  202. { type: 'step/end', seq: 5, time: 6, data: { turn: 1, step: 1 } },
  203. { type: 'turn/end', seq: 6, time: 7, data: { turn: 1, reason: { kind: 'completed' } } },
  204. ]
  205. await ctx.sessionPersistence.create(m)
  206. await ctx.sessionPersistence.append(m.id, log)
  207. const loaded = await ctx.sessionPersistence.load(m.id)
  208. expect(loaded.events).toEqual(log) // chunks preserved, contiguous seqs
  209. })
  210. it('source-qualifies revisions across roots while preserving same-log reopen identity', async () => {
  211. const m = meta('revision-source')
  212. await ctx.sessionPersistence.create(m)
  213. await ctx.sessionPersistence.append(m.id, oneTurnLog())
  214. const revision = (await ctx.sessionPersistence.listSnapshots())[0]?.revision
  215. const reopenedCtx = new Context()
  216. await reopenedCtx.plugin(SessionStore)
  217. await reopenedCtx.plugin(SessionPersistenceJsonl, { root, compression: 'none' })
  218. expect((await reopenedCtx.sessionPersistence.listSnapshots())[0]?.revision).toBe(revision)
  219. const otherRoot = await freshRoot()
  220. const otherCtx = new Context()
  221. await otherCtx.plugin(SessionStore)
  222. await otherCtx.plugin(SessionPersistenceJsonl, { root: otherRoot, compression: 'none' })
  223. await otherCtx.sessionPersistence.create(m)
  224. await otherCtx.sessionPersistence.append(m.id, oneTurnLog())
  225. expect((await otherCtx.sessionPersistence.listSnapshots())[0]?.revision).not.toBe(revision)
  226. await reopenedCtx.fiber.dispose()
  227. await otherCtx.fiber.dispose()
  228. })
  229. it('omits a snapshot artifact removed after discovery', async () => {
  230. const m = meta('vanishing-snapshot')
  231. await ctx.sessionPersistence.create(m)
  232. await ctx.sessionPersistence.append(m.id, oneTurnLog())
  233. const persistence = ctx.sessionPersistence as unknown as {
  234. listArtifacts(): Promise<Array<{ header: SessionHeader; path: string }>>
  235. }
  236. const listArtifacts = persistence.listArtifacts.bind(persistence)
  237. const discovery = vi.spyOn(persistence, 'listArtifacts').mockImplementation(async () => {
  238. const artifacts = await listArtifacts()
  239. await rm(artifacts[0]!.path)
  240. return artifacts
  241. })
  242. await expect(ctx.sessionPersistence.listSnapshots()).resolves.toEqual([])
  243. discovery.mockRestore()
  244. })
  245. it('surfaces non-ENOENT snapshot stat failures after discovery', async () => {
  246. const persistence = ctx.sessionPersistence as unknown as {
  247. listArtifacts(): Promise<Array<{ header: SessionHeader; path: string }>>
  248. }
  249. const discovery = vi.spyOn(persistence, 'listArtifacts').mockResolvedValue([{
  250. header: meta('snapshot-stat-failure'),
  251. path: `${root}\0snapshot-stat-failure`,
  252. }])
  253. await expect(ctx.sessionPersistence.listSnapshots()).rejects.toThrow(/null bytes/)
  254. discovery.mockRestore()
  255. })
  256. it('forwards snapshot-list cancellation and awaits in-flight discovery cleanup', async () => {
  257. const persistence = ctx.sessionPersistence as unknown as {
  258. listArtifacts(signal?: AbortSignal): Promise<Array<{ header: SessionHeader; path: string }>>
  259. }
  260. const started = Promise.withResolvers<AbortSignal>()
  261. const cleanup = Promise.withResolvers<undefined>()
  262. vi.spyOn(persistence, 'listArtifacts').mockImplementation(async (signal) => {
  263. if (signal === undefined) throw new Error('expected snapshot-list signal')
  264. started.resolve(signal)
  265. await cleanup.promise
  266. return []
  267. })
  268. const reason = new Error('JSONL snapshot discovery cancelled')
  269. const controller = new AbortController()
  270. const pending = ctx.sessionPersistence.listSnapshots(controller.signal)
  271. expect(await started.promise).toBe(controller.signal)
  272. let settled = false
  273. void pending.then(
  274. () => { settled = true },
  275. () => { settled = true },
  276. )
  277. controller.abort(reason)
  278. await Promise.resolve()
  279. expect(settled).toBe(false)
  280. cleanup.resolve(undefined)
  281. await expect(pending).rejects.toBe(reason)
  282. })
  283. it('checks cancellation after an uncancellable snapshot stat settles', async () => {
  284. const m = meta('snapshot-stat-cancellation')
  285. await ctx.sessionPersistence.create(m)
  286. await ctx.sessionPersistence.append(m.id, oneTurnLog())
  287. const persistence = ctx.sessionPersistence as unknown as {
  288. listArtifacts(signal?: AbortSignal): Promise<Array<{ header: SessionHeader; path: string }>>
  289. }
  290. const discovery = vi.spyOn(persistence, 'listArtifacts').mockResolvedValue([{
  291. header: m,
  292. path: rawLogPath(root, m.cwd, m.id),
  293. }])
  294. const reason = new Error('JSONL snapshot stat cancelled')
  295. const controller = new AbortController()
  296. const pending = ctx.sessionPersistence.listSnapshots(controller.signal)
  297. queueMicrotask(() => { controller.abort(reason) })
  298. await expect(pending).rejects.toBe(reason)
  299. expect(discovery).toHaveBeenCalledWith(controller.signal)
  300. })
  301. it('rejects a stored v0 log containing a legacy request/header-delta event', async () => {
  302. const m = meta('legacy-header-delta', '/legacy')
  303. const path = rawLogPath(root, m.cwd, m.id)
  304. await mkdir(sessionDir(root, m.cwd, m.id), { recursive: true })
  305. await writeFile(path, [
  306. JSON.stringify(toHeaderLine(m)),
  307. JSON.stringify({ type: 'turn/start', seq: 0, time: 1, data: { turn: 1 } }),
  308. JSON.stringify({ type: 'request/header-delta', seq: 1, time: 2, data: { config: { model: 'legacy' } } }),
  309. JSON.stringify({ type: 'turn/end', seq: 2, time: 3, data: { turn: 1, reason: { kind: 'completed' } } }),
  310. '',
  311. ].join('\n'))
  312. await expect(ctx.sessionPersistence.load(m.id)).rejects.toThrow(/unsupported legacy request\/header-delta event at seq 1/)
  313. })
  314. it('rejects a stored v0 full header carrying the legacy fallback reason', async () => {
  315. const m = meta('legacy-header-fallback', '/legacy')
  316. const path = rawLogPath(root, m.cwd, m.id)
  317. await mkdir(sessionDir(root, m.cwd, m.id), { recursive: true })
  318. await writeFile(path, [
  319. JSON.stringify(toHeaderLine(m)),
  320. JSON.stringify({
  321. type: 'request/header',
  322. seq: 0,
  323. time: 1,
  324. data: { header: { config: { model: 'legacy' } }, reason: 'fallback' },
  325. }),
  326. '',
  327. ].join('\n'))
  328. await expect(ctx.sessionPersistence.load(m.id))
  329. .rejects.toThrow(/unsupported legacy request\/header reason "fallback" at seq 0/)
  330. })
  331. it('persists a forked child seed through the existing session write path', async () => {
  332. const source = ctx.sessions.create(SessionId('persist-parent'), { meta: { cwd: '/workspace' } })
  333. appendClosedTurn(source)
  334. const child = ctx.sessions.fork(source, undefined, SessionId('persist-child'))
  335. await ctx.sessions.flush(child)
  336. const loaded = await ctx.sessionPersistence.load(child.id)
  337. // The constructor seed reaches disk verbatim, then the child's end-seed.
  338. expect(loaded.events.slice(0, source.events.length)).toEqual(source.events)
  339. expect(loaded.events.at(-1)).toMatchObject({ type: 'session/end-seed', seq: source.events.length })
  340. expect(loaded.meta).toMatchObject({
  341. id: SessionId('persist-child'),
  342. cwd: '/workspace',
  343. parentSession: SessionId('persist-parent'),
  344. seedLength: source.events.length,
  345. })
  346. })
  347. it('crash recovery: load preserves the interrupted turn and closes it with a synthetic turn/end {interrupted}', async () => {
  348. const m = meta('crash', '/proj')
  349. await ctx.sessionPersistence.create(m)
  350. await ctx.sessionPersistence.append(m.id, oneTurnLog()) // seqs 0..5, turn/end at 5
  351. // Simulate a crash mid-second-turn: append raw lines that are NOT closed by
  352. // a turn/end (turn/start + step/start are fully written), plus a final
  353. // partial line with no newline (a torn fragment never fully flushed).
  354. const path = rawLogPath(root, '/proj', m.id)
  355. await writeFile(path, [
  356. JSON.stringify({ type: 'turn/start', seq: 6, time: 8, data: { turn: 2 } }),
  357. JSON.stringify({ type: 'step/start', seq: 7, time: 9, data: { turn: 2, step: 1 } }),
  358. '{"type":"assistant/chunk","seq":8,"ti', // truncated partial line (no newline)
  359. ].join('\n'), { flag: 'a' })
  360. // load PRESERVES the interrupted turn's real events (turn/start 6, step/start
  361. // 7) — a turn can be huge, so they must not be truncated — and durably closes
  362. // the orphaned turn with synthetic step/end (8) + turn/end {interrupted} (9).
  363. const loaded = await ctx.sessionPersistence.load(m.id)
  364. expect(loaded.events.map(e => e.seq)).toEqual([0, 1, 2, 3, 4, 5, 6, 7, 8, 9])
  365. const last = loaded.events.at(-1)!
  366. expect(last.type === 'turn/end' && last.data.reason).toEqual({ kind: 'interrupted' })
  367. const stepEnd = loaded.events[8]!
  368. expect(stepEnd.type).toBe('step/end')
  369. // the torn seq-8 chunk fragment did not survive
  370. expect(loaded.events.some(e => e.type === 'assistant/chunk' && e.seq === 8)).toBe(false)
  371. // The next append continues at seq 10 (the balanced length).
  372. const turn3 = [
  373. { type: 'turn/start', seq: 10, time: 11, data: { turn: 3 } },
  374. { type: 'turn/end', seq: 11, time: 12, data: { turn: 3, reason: { kind: 'completed' } } },
  375. ] as SessionEvent[]
  376. await ctx.sessionPersistence.append(m.id, turn3)
  377. const reloaded = await ctx.sessionPersistence.load(m.id)
  378. expect(reloaded.events.map(e => e.seq)).toEqual([0, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11])
  379. })
  380. it('committed events are never rewritten: only the crash tail is repaired', async () => {
  381. const m = meta('append-only')
  382. await ctx.sessionPersistence.create(m)
  383. await ctx.sessionPersistence.append(m.id, oneTurnLog())
  384. const before = await readFile(rawLogPath(root, undefined, m.id), 'utf8')
  385. const committedPrefix = before // the whole committed log
  386. // A crash tail then a repair-append.
  387. await writeFile(rawLogPath(root, undefined, m.id), '\n{"partial', { flag: 'a' })
  388. await ctx.sessionPersistence.load(m.id)
  389. await ctx.sessionPersistence.append(m.id, [
  390. { type: 'turn/start', seq: 6, time: 9, data: { turn: 2 } },
  391. { type: 'turn/end', seq: 7, time: 10, data: { turn: 2, reason: { kind: 'completed' } } },
  392. ] as SessionEvent[])
  393. const after = await readFile(rawLogPath(root, undefined, m.id), 'utf8')
  394. // the committed prefix is byte-for-byte intact at the head of the file
  395. expect(after.startsWith(committedPrefix)).toBe(true)
  396. })
  397. it('a failed appendLines truncates partial bytes so a retry has no seq gap', async () => {
  398. const m = meta('truncate-retry')
  399. await ctx.sessionPersistence.create(m)
  400. await ctx.sessionPersistence.append(m.id, oneTurnLog()) // materialized, seqs 0..5
  401. const sizeBefore = (await stat(rawLogPath(root, undefined, m.id))).size
  402. // Force the NEXT fsync (inside appendLines) to fail once, AFTER writeFile
  403. // has already put bytes on disk — simulating an ENOSPC/fsync error
  404. // mid-append. The recovery truncate() also fsyncs, so allow that one.
  405. const handle = await (await import('node:fs/promises')).open(rawLogPath(root, undefined, m.id), 'r')
  406. const proto = Object.getPrototypeOf(handle) as { sync: () => Promise<void> }
  407. await handle.close()
  408. const realSync = proto.sync
  409. let failed = false
  410. const spy = vi.spyOn(proto, 'sync').mockImplementation(async function (this: unknown) {
  411. if (!failed) { failed = true; throw new Error('simulated fsync ENOSPC') }
  412. return realSync.call(this)
  413. })
  414. const turn2 = [
  415. { type: 'turn/start', seq: 6, time: 9, data: { turn: 2 } },
  416. { type: 'turn/end', seq: 7, time: 10, data: { turn: 2, reason: { kind: 'completed' } } },
  417. ] as SessionEvent[]
  418. // The append rejects, but the partial bytes are truncated back: the file is
  419. // its pre-append size and the cursor is unchanged.
  420. await expect(ctx.sessionPersistence.append(m.id, turn2)).rejects.toThrow(/ENOSPC/)
  421. expect((await stat(rawLogPath(root, undefined, m.id))).size).toBe(sizeBefore)
  422. spy.mockRestore()
  423. // The retry now succeeds with NO seq gap — the log is contiguous 0..7.
  424. await ctx.sessionPersistence.append(m.id, turn2)
  425. const loaded = await ctx.sessionPersistence.load(m.id)
  426. expect(loaded.events.map(e => e.seq)).toEqual([0, 1, 2, 3, 4, 5, 6, 7])
  427. })
  428. it('reports both the append failure and a failed rollback', async () => {
  429. const m = meta('rollback-failure')
  430. await ctx.sessionPersistence.create(m)
  431. await ctx.sessionPersistence.append(m.id, oneTurnLog())
  432. const path = rawLogPath(root, undefined, m.id)
  433. const handle = await (await import('node:fs/promises')).open(path, 'r')
  434. const proto = Object.getPrototypeOf(handle) as { sync: () => Promise<void> }
  435. await handle.close()
  436. const realSync = proto.sync
  437. let failed = false
  438. const syncSpy = vi.spyOn(proto, 'sync').mockImplementation(async function (this: unknown) {
  439. if (!failed) { failed = true; throw new Error('simulated append fsync failure') }
  440. return realSync.call(this)
  441. })
  442. const backend = ctx.sessionPersistence as unknown as {
  443. rollbackAppend: (path: string, size: number) => Promise<void>
  444. }
  445. const realRollback = backend.rollbackAppend.bind(backend)
  446. backend.rollbackAppend = () => Promise.reject(new Error('simulated rollback failure'))
  447. try {
  448. await ctx.sessionPersistence.append(m.id, [
  449. { type: 'turn/start', seq: 6, time: 9, data: { turn: 2 } },
  450. ] as SessionEvent[])
  451. throw new Error('expected append to reject')
  452. } catch (error) {
  453. expect(error).toBeInstanceOf(AggregateError)
  454. const aggregate = error as AggregateError
  455. expect(aggregate.message).toContain(`failed to roll back append to "${path}"`)
  456. expect(aggregate.errors).toHaveLength(2)
  457. expect(aggregate.errors[0]).toMatchObject({ message: 'simulated append fsync failure' })
  458. expect(aggregate.errors[1]).toMatchObject({ message: 'simulated rollback failure' })
  459. } finally {
  460. backend.rollbackAppend = realRollback
  461. syncSpy.mockRestore()
  462. }
  463. })
  464. it('load returns a meta copy: mutating it does not corrupt backend pathing', async () => {
  465. const m = meta('meta-copy', '/proj')
  466. await ctx.sessionPersistence.create(m)
  467. await ctx.sessionPersistence.append(m.id, oneTurnLog())
  468. const loaded = await ctx.sessionPersistence.load(m.id)
  469. // A consumer mutates the returned meta's cwd. The backend's stored pathing
  470. // metadata must be unaffected, so a later append still finds the right log.
  471. mutableHeader(loaded.meta).cwd = '/evil'
  472. await ctx.sessionPersistence.append(m.id, [
  473. { type: 'turn/start', seq: 6, time: 9, data: { turn: 2 } },
  474. { type: 'turn/end', seq: 7, time: 10, data: { turn: 2, reason: { kind: 'completed' } } },
  475. ] as SessionEvent[])
  476. // The append landed in the ORIGINAL /proj log, not beside an /evil path.
  477. const reloaded = await ctx.sessionPersistence.load(m.id)
  478. expect(reloaded.meta.cwd).toBe('/proj')
  479. expect(reloaded.events.map(e => e.seq)).toEqual([0, 1, 2, 3, 4, 5, 6, 7])
  480. })
  481. it('rejects a mismatched header before repairing either session log', async () => {
  482. const a = meta('identity-a', '/same')
  483. const b = meta('identity-b', '/same')
  484. await ctx.sessionPersistence.create(a)
  485. await ctx.sessionPersistence.append(a.id, [{
  486. type: 'turn/start',
  487. seq: 0,
  488. time: 1,
  489. data: { turn: 1 },
  490. }])
  491. await ctx.sessionPersistence.create(b)
  492. await ctx.sessionPersistence.append(b.id, oneTurnLog())
  493. const aPath = rawLogPath(root, a.cwd, a.id)
  494. const bPath = rawLogPath(root, b.cwd, b.id)
  495. await rewriteHeader(aPath, (header) => { header.id = b.id })
  496. const beforeA = await readFile(aPath)
  497. const beforeB = await readFile(bPath)
  498. await expect(ctx.sessionPersistence.load(a.id))
  499. .rejects.toThrow(/requested id "identity-a" does not match header id "identity-b"/)
  500. expect(await readFile(aPath)).toEqual(beforeA)
  501. expect(await readFile(bPath)).toEqual(beforeB)
  502. })
  503. it('rejects a re-append of an already-stored seq', async () => {
  504. const m = meta('reappend')
  505. await ctx.sessionPersistence.create(m)
  506. await ctx.sessionPersistence.append(m.id, oneTurnLog())
  507. await expect(ctx.sessionPersistence.append(m.id, oneTurnLog())).rejects.toThrow(/seq mismatch/)
  508. })
  509. it('path-traversal session ids are neutralized (no escape from root)', async () => {
  510. const evil = SessionId('../../etc/pwn')
  511. const m = { version: 0, id: evil, createdAt: 1 }
  512. await ctx.sessionPersistence.create(m)
  513. await ctx.sessionPersistence.append(evil, oneTurnLog())
  514. // The file lives UNDER root, not at ../../etc.
  515. const all: string[] = []
  516. async function walk(dir: string): Promise<void> {
  517. for (const e of await readdir(dir, { withFileTypes: true })) {
  518. const p = join(dir, e.name)
  519. if (e.isDirectory()) await walk(p)
  520. else all.push(p)
  521. }
  522. }
  523. await walk(root)
  524. expect(all.length).toBeGreaterThan(0)
  525. expect(all.every(p => p.startsWith(root))).toBe(true)
  526. })
  527. })
  528. describe('SessionPersistenceJsonl: write path (session/event → flush)', () => {
  529. it('concurrent sessions do not cross buffers', async () => {
  530. root = await freshRoot()
  531. const ctx = new Context()
  532. await ctx.plugin(SessionStore)
  533. await ctx.plugin(SessionPersistenceJsonl, { root, compression: 'none' })
  534. const a = ctx.sessions.create(SessionId('sa'))
  535. const b = ctx.sessions.create(SessionId('sb'))
  536. a.append('turn/start', { turn: 1 })
  537. b.append('turn/start', { turn: 1 })
  538. a.append('user/message', createUserMessage({
  539. content: [{ type: 'text', text: 'A' }], source: { kind: 'user' },
  540. }), { surfaceOp: 'append' })
  541. b.append('user/message', createUserMessage({
  542. content: [{ type: 'text', text: 'B' }], source: { kind: 'user' },
  543. }), { surfaceOp: 'append' })
  544. a.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  545. b.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  546. await ctx.sessions.flush(a)
  547. await ctx.sessions.flush(b)
  548. const la = await ctx.sessionPersistence.load(SessionId('sa'))
  549. const lb = await ctx.sessionPersistence.load(SessionId('sb'))
  550. expect(JSON.stringify(la.events)).toContain('"A"')
  551. expect(JSON.stringify(la.events)).not.toContain('"B"')
  552. expect(JSON.stringify(lb.events)).toContain('"B"')
  553. expect(JSON.stringify(lb.events)).not.toContain('"A"')
  554. await ctx.fiber.dispose()
  555. })
  556. })
  557. describe('SessionPersistenceJsonl: scanLog unit', () => {
  558. it('rejects a header-less / empty log', () => {
  559. expect(() => scanLog(Buffer.from(''))).toThrow()
  560. })
  561. it('rejects a corrupt header line', () => {
  562. expect(() => scanLog(Buffer.from('not json\n'))).toThrow(/header/)
  563. })
  564. it('rejects a non-session first line', () => {
  565. expect(() => scanLog(Buffer.from('{"type":"event"}\n'))).toThrow(/session header/)
  566. })
  567. it.each([
  568. ['fractional', 1.5],
  569. ['negative', -1],
  570. ['unsafe', Number.MAX_SAFE_INTEGER + 1],
  571. ])('rejects a session header with a %s createdAt', (_label, createdAt) => {
  572. const log = JSON.stringify({
  573. type: 'session',
  574. version: 0,
  575. id: 'invalid-created-at',
  576. createdAt,
  577. delegationDepth: 0,
  578. }) + '\n'
  579. expect(() => scanLog(Buffer.from(log))).toThrow(/session header/)
  580. })
  581. it('rejects a session header with negative-zero createdAt', () => {
  582. const log = '{"type":"session","version":0,"id":"invalid-created-at","createdAt":-0,"delegationDepth":0}\n'
  583. expect(() => scanLog(Buffer.from(log))).toThrow(/session header/)
  584. })
  585. it.each([
  586. ['missing', undefined],
  587. ['a string', '1'],
  588. ['fractional', 1.5],
  589. ['negative', -1],
  590. ])('rejects a session header with %s delegationDepth', (_label, delegationDepth) => {
  591. const log = JSON.stringify({
  592. type: 'session',
  593. version: 0,
  594. id: 'invalid-depth',
  595. createdAt: 1,
  596. ...delegationDepth === undefined ? {} : { delegationDepth },
  597. }) + '\n'
  598. expect(() => scanLog(Buffer.from(log))).toThrow(/session header/)
  599. })
  600. it('rejects a session header with negative-zero delegationDepth', () => {
  601. const log = '{"type":"session","version":0,"id":"invalid-depth","createdAt":1,"delegationDepth":-0}\n'
  602. expect(() => scanLog(Buffer.from(log))).toThrow(/session header/)
  603. })
  604. it('a seq gap after the last turn/end bounds the preserved tail (torn fragment tolerated)', () => {
  605. const log = [
  606. JSON.stringify({ type: 'session', version: 0, id: 'g', createdAt: 1, delegationDepth: 0 }),
  607. JSON.stringify({ type: 'turn/start', seq: 0, time: 1, data: { turn: 1 } }),
  608. JSON.stringify({ type: 'step/start', seq: 2, time: 2, data: { turn: 1, step: 1 } }), // gap: missing seq 1
  609. ].join('\n') + '\n'
  610. // No committed turn/end, so the gap is a tolerated crash boundary: scanLog PRESERVES the
  611. // contiguous prefix (turn/start seq 0) — real interrupted-turn work, not discarded — and
  612. // stops at the gap. `loadCore`, not this scanner, later closes the orphaned turn.
  613. expect(scanLog(Buffer.from(log)).events.map(e => e.seq)).toEqual([0])
  614. })
  615. it('rejects a seq gap BEFORE a later committed turn/end (committed data damaged)', () => {
  616. const log = [
  617. JSON.stringify({ type: 'session', version: 0, id: 'g2', createdAt: 1, delegationDepth: 0 }),
  618. JSON.stringify({ type: 'turn/start', seq: 0, time: 1, data: { turn: 1 } }),
  619. JSON.stringify({ type: 'step/start', seq: 2, time: 2, data: { turn: 1, step: 1 } }), // gap: missing seq 1
  620. JSON.stringify({ type: 'turn/end', seq: 3, time: 3, data: { turn: 1, reason: { kind: 'completed' } } }),
  621. ].join('\n') + '\n'
  622. // A turn/end exists, so the prefix up to it is committed — but it has a hole.
  623. // Truncating it would silently drop committed data → unloadable.
  624. expect(() => scanLog(Buffer.from(log))).toThrow(/seq gap in committed region/)
  625. })
  626. it('rejects a corrupt line BEFORE a later committed turn/end (committed data damaged)', () => {
  627. const log = [
  628. JSON.stringify({ type: 'session', version: 0, id: 'c', createdAt: 1, delegationDepth: 0 }),
  629. '{not json', // corrupt, sits in the committed region (a turn/end follows)
  630. JSON.stringify({ type: 'turn/end', seq: 1, time: 2, data: { turn: 1, reason: { kind: 'completed' } } }),
  631. ].join('\n') + '\n'
  632. expect(() => scanLog(Buffer.from(log))).toThrow(/unparsable committed event/)
  633. })
  634. it('a header-only log (no event lines at all) preserves nothing — committedBytes is the header', () => {
  635. const log = JSON.stringify({ type: 'session', version: 0, id: 'h0', createdAt: 1, delegationDepth: 0 }) + '\n'
  636. const scanned = scanLog(Buffer.from(log))
  637. expect(scanned.events).toEqual([])
  638. // committedBytes falls back to the header line's end (no preserved events).
  639. expect(scanned.committedBytes).toBe(Buffer.byteLength(log, 'utf8'))
  640. })
  641. it('a corrupt line after the last turn/end bounds the preserved tail', () => {
  642. const log = [
  643. JSON.stringify({ type: 'session', version: 0, id: 'c2', createdAt: 1, delegationDepth: 0 }),
  644. JSON.stringify({ type: 'turn/start', seq: 0, time: 1, data: { turn: 1 } }),
  645. '{not json', // corrupt crash fragment, no turn/end committed
  646. ].join('\n') + '\n'
  647. // The contiguous prefix (turn/start seq 0) is preserved; the corrupt
  648. // fragment after it is the tolerated crash boundary.
  649. expect(scanLog(Buffer.from(log)).events.map(e => e.seq)).toEqual([0])
  650. })
  651. it('tolerates a seq gap AFTER a turn/end (uncommitted tail)', () => {
  652. const log = [
  653. JSON.stringify({ type: 'session', version: 0, id: 't', createdAt: 1, delegationDepth: 0 }),
  654. JSON.stringify({ type: 'turn/start', seq: 0, time: 1, data: { turn: 1 } }),
  655. JSON.stringify({ type: 'turn/end', seq: 1, time: 2, data: { turn: 1, reason: { kind: 'completed' } } }),
  656. JSON.stringify({ type: 'step/start', seq: 9, time: 3, data: { turn: 2, step: 1 } }), // gap in uncommitted tail
  657. ].join('\n') + '\n'
  658. const { events } = scanLog(Buffer.from(log))
  659. expect(events.map(e => e.seq)).toEqual([0, 1]) // tail dropped
  660. })
  661. })
  662. describe('SessionPersistenceJsonl: default packed chunk rows', () => {
  663. let ctx: Context
  664. beforeEach(async () => {
  665. root = await freshRoot()
  666. ctx = new Context()
  667. await ctx.plugin(SessionStore)
  668. // compression: 'none' — these tests assert the textual storage-record layout
  669. // (row tags per line); packing is orthogonal to the physical encoding.
  670. await ctx.plugin(SessionPersistenceJsonl, { root, compression: 'none' })
  671. })
  672. afterEach(async () => { await ctx.fiber.dispose() })
  673. /** A one-turn log whose step streams a five-member text-delta run. */
  674. function chunkRunLog(): SessionEvent[] {
  675. const deltas: SessionEvent[] = Array.from({ length: 5 }, (_, k) => ({
  676. type: 'assistant/chunk',
  677. seq: 2 + k,
  678. time: 3 + k,
  679. data: { turn: 1, step: 1, chunk: { type: 'text-delta', index: 0, text: `t${k}` } },
  680. }))
  681. return [
  682. { type: 'turn/start', seq: 0, time: 1, data: { turn: 1 } },
  683. { type: 'step/start', seq: 1, time: 2, data: { turn: 1, step: 1 } },
  684. ...deltas,
  685. { type: 'assistant/message', seq: 7, time: 8, data: {
  686. turn: 1, step: 1,
  687. message: createMessage({
  688. role: 'assistant',
  689. content: [{ type: 'text', text: 't0t1t2t3t4' }],
  690. source: {
  691. kind: 'model',
  692. ...{ provider: 'mock', model: 'mock' },
  693. },
  694. }),
  695. }, surfaceOp: 'append', sourceEventSeqs: [2, 3, 4, 5, 6] },
  696. { type: 'step/end', seq: 8, time: 9, data: { turn: 1, step: 1 } },
  697. { type: 'turn/end', seq: 9, time: 10, data: { turn: 1, reason: { kind: 'completed' } } },
  698. ]
  699. }
  700. it('writes a delta run as one text-chunks row by default and loads back identical events', async () => {
  701. const m = meta('packed', '/work')
  702. const log = chunkRunLog()
  703. await ctx.sessionPersistence.create(m)
  704. await ctx.sessionPersistence.append(m.id, log)
  705. const raw = (await readFile(rawLogPath(root, '/work', m.id), 'utf8')).split('\n').filter(Boolean)
  706. const tags = raw.slice(1).map(line => (JSON.parse(line) as { type: string }).type)
  707. expect(tags).toEqual(['turn/start', 'step/start', 'text-chunks', 'assistant/message', 'step/end', 'turn/end'])
  708. const loaded = await ctx.sessionPersistence.load(m.id)
  709. expect(loaded.events).toEqual(log)
  710. })
  711. it('packChunks: false writes one event per line and still loads identical events', async () => {
  712. const unpackedRoot = await freshRoot()
  713. const unpacked = new Context()
  714. await unpacked.plugin(SessionStore)
  715. await unpacked.plugin(SessionPersistenceJsonl, {
  716. root: unpackedRoot,
  717. packChunks: false,
  718. compression: 'none',
  719. })
  720. try {
  721. const m = meta('unpacked', '/work')
  722. const log = chunkRunLog()
  723. await unpacked.sessionPersistence.create(m)
  724. await unpacked.sessionPersistence.append(m.id, log)
  725. const records = (await readFile(rawLogPath(unpackedRoot, '/work', m.id), 'utf8'))
  726. .split('\n').filter(Boolean).slice(1)
  727. .map(line => JSON.parse(line) as { type: string })
  728. expect(records.filter(record => record.type === 'assistant/chunk')).toHaveLength(5)
  729. expect(records.some(record => record.type === 'text-chunks')).toBe(false)
  730. expect((await unpacked.sessionPersistence.load(m.id)).events).toEqual(log)
  731. } finally {
  732. await unpacked.fiber.dispose()
  733. }
  734. })
  735. it('loads a mixed file: verbatim lines from an unpacked writer, then packed appends', async () => {
  736. const m = meta('mixed', '/work')
  737. const log = chunkRunLog()
  738. // First turn written line-per-event by an unpacked-config writer (an old
  739. // file, hand-planted so this packed-config backend adopts it on load).
  740. await mkdir(sessionDir(root, '/work', m.id), { recursive: true })
  741. await writeFile(rawLogPath(root, '/work', m.id), [
  742. JSON.stringify({ type: 'session', version: 0, id: 'mixed', createdAt: 1000, cwd: '/work', delegationDepth: 0 }),
  743. ...log.map(e => JSON.stringify(e)),
  744. ].join('\n') + '\n')
  745. // Adopt the stored log (cursor = stored length), then append a second turn
  746. // through THIS packed-config backend.
  747. expect((await ctx.sessionPersistence.load(m.id)).events).toEqual(log)
  748. const secondTurn: SessionEvent[] = JSON.parse(JSON.stringify(log)) as SessionEvent[]
  749. for (const [k, e] of secondTurn.entries()) {
  750. ;(e as { seq: number }).seq = 10 + k
  751. ;(e.data as { turn: number }).turn = 2
  752. }
  753. await ctx.sessionPersistence.append(m.id, secondTurn)
  754. const loaded = await ctx.sessionPersistence.load(m.id)
  755. expect(loaded.events).toEqual([...log, ...secondTurn])
  756. // The packed append really packed: the file's tail carries a text-chunks row.
  757. const tags = (await readFile(rawLogPath(root, '/work', m.id), 'utf8')).split('\n').filter(Boolean)
  758. .map(line => (JSON.parse(line) as { type: string }).type)
  759. expect(tags.filter(t => t === 'text-chunks')).toHaveLength(1)
  760. expect(tags.filter(t => t === 'assistant/chunk')).toHaveLength(5)
  761. })
  762. it('scanLog: a packed row advances the seq cursor by its whole run', () => {
  763. const logText = [
  764. JSON.stringify({ type: 'session', version: 0, id: 'rows', createdAt: 1, delegationDepth: 0 }),
  765. JSON.stringify({ type: 'turn/start', seq: 0, time: 1, data: { turn: 1 } }),
  766. JSON.stringify({ type: 'text-chunks', seq0: 1, time0: 2, data: { turn: 1, step: 1, index: 0, dt: [1, 1], texts: ['a', 'b', 'c'] } }),
  767. JSON.stringify({ type: 'turn/end', seq: 4, time: 5, data: { turn: 1, reason: { kind: 'completed' } } }),
  768. ].join('\n') + '\n'
  769. const { events } = scanLog(Buffer.from(logText))
  770. expect(events.map(e => e.seq)).toEqual([0, 1, 2, 3, 4])
  771. expect(events[2]).toEqual({ type: 'assistant/chunk', seq: 2, time: 3, data: { turn: 1, step: 1, chunk: { type: 'text-delta', index: 0, text: 'b' } } })
  772. })
  773. it('scanLog: a malformed packed row in the committed region rejects like corrupt JSON', () => {
  774. const logText = [
  775. JSON.stringify({ type: 'session', version: 0, id: 'bad-row', createdAt: 1, delegationDepth: 0 }),
  776. // dt arity mismatch — row validation throws, so the line is a committed hole.
  777. JSON.stringify({ type: 'text-chunks', seq0: 0, time0: 1, data: { turn: 1, step: 1, index: 0, dt: [], texts: ['a', 'b'] } }),
  778. JSON.stringify({ type: 'turn/end', seq: 2, time: 3, data: { turn: 1, reason: { kind: 'completed' } } }),
  779. ].join('\n') + '\n'
  780. expect(() => scanLog(Buffer.from(logText))).toThrow(/unparsable committed event/)
  781. })
  782. it('scanLog: a packed row with a mid-run seq gap after the last turn/end drops the whole row', () => {
  783. const logText = [
  784. JSON.stringify({ type: 'session', version: 0, id: 'row-gap', createdAt: 1, delegationDepth: 0 }),
  785. JSON.stringify({ type: 'turn/start', seq: 0, time: 1, data: { turn: 1 } }),
  786. // seq0 skips 1 — the run's first member is already a gap; no turn/end follows.
  787. JSON.stringify({ type: 'text-chunks', seq0: 2, time0: 2, data: { turn: 1, step: 1, index: 0, dt: [1, 1], texts: ['a', 'b', 'c'] } }),
  788. ].join('\n') + '\n'
  789. const scanned = scanLog(Buffer.from(logText))
  790. expect(scanned.events.map(e => e.seq)).toEqual([0])
  791. // committedBytes stays on the line boundary BEFORE the dropped row.
  792. const headerAndTurn = logText.split('\n').slice(0, 2).join('\n') + '\n'
  793. expect(scanned.committedBytes).toBe(Buffer.byteLength(headerAndTurn, 'utf8'))
  794. })
  795. it('eventLines(packChunks: false) is byte-identical to the pre-packing layout', () => {
  796. const log = chunkRunLog()
  797. expect(eventLines(log, false)).toBe(log.map(e => JSON.stringify(e)).join('\n'))
  798. })
  799. })
  800. describe('SessionPersistenceJsonl: edge cases', () => {
  801. let ctx: Context
  802. beforeEach(async () => {
  803. root = await freshRoot()
  804. ctx = new Context()
  805. await ctx.plugin(SessionStore)
  806. await ctx.plugin(SessionPersistenceJsonl, { root, compression: 'none' })
  807. })
  808. afterEach(async () => { await ctx.fiber.dispose() })
  809. it('append rejects non-JSON-serializable undefined-producing data', async () => {
  810. const m = meta('undef')
  811. await ctx.sessionPersistence.create(m)
  812. // A value whose JSON.stringify yields undefined (a bare function as data).
  813. const bad = [{ type: 'user/message', seq: 0, time: 1, data: (() => 0) as unknown }] as unknown as SessionEvent[]
  814. await expect(ctx.sessionPersistence.append(m.id, bad)).rejects.toThrow(/non-JSON-serializable/)
  815. })
  816. it('create snapshots its meta: mutating the caller object after the call is ignored', async () => {
  817. const m = meta('create-snap', '/orig')
  818. const p = ctx.sessionPersistence.create(m)
  819. // Mutate the caller's meta object immediately after calling create.
  820. mutableHeader(m).cwd = '/mutated'
  821. await p
  822. await ctx.sessionPersistence.append(SessionId('create-snap'), oneTurnLog())
  823. // The log materialized under the ORIGINAL cwd, not the mutated one.
  824. expect((await stat(rawLogPath(root, '/orig', SessionId('create-snap')))).isFile()).toBe(true)
  825. await expect(stat(rawLogPath(root, '/mutated', SessionId('create-snap')))).rejects.toThrow()
  826. })
  827. it('list discovers sessions across multiple project directories', async () => {
  828. await ctx.sessionPersistence.create(meta('p1', '/projA'))
  829. await ctx.sessionPersistence.append(SessionId('p1'), oneTurnLog())
  830. await ctx.sessionPersistence.create(meta('p2', '/projB'))
  831. await ctx.sessionPersistence.append(SessionId('p2'), oneTurnLog())
  832. await ctx.sessionPersistence.create(meta('p3')) // no cwd → _no-cwd project directory
  833. await ctx.sessionPersistence.append(SessionId('p3'), oneTurnLog())
  834. const ids = (await ctx.sessionPersistence.list()).map(x => x.id).sort()
  835. expect(ids).toEqual(['p1', 'p2', 'p3'])
  836. })
  837. it('groups sessions whose cwd paths normalize to the same project directory', async () => {
  838. const first = meta('normalized-first', '/a/b-c')
  839. const second = meta('normalized-second', '/a-b/c')
  840. await ctx.sessionPersistence.create(first)
  841. await ctx.sessionPersistence.append(first.id, oneTurnLog())
  842. await ctx.sessionPersistence.create(second)
  843. await ctx.sessionPersistence.append(second.id, oneTurnLog())
  844. expect(projectDir(root, first.cwd)).toBe(projectDir(root, second.cwd))
  845. expect(await readdir(projectDir(root, first.cwd))).toEqual(expect.arrayContaining([
  846. encodeSegment(first.id),
  847. encodeSegment(second.id),
  848. ]))
  849. expect((await ctx.sessionPersistence.list()).map(header => header.id).sort())
  850. .toEqual([first.id, second.id].sort())
  851. })
  852. it('list on an empty root returns nothing', async () => {
  853. expect(await ctx.sessionPersistence.list()).toEqual([])
  854. })
  855. it('keeps the transcript in an extensible session-owned directory', async () => {
  856. const m = meta('owned-directory', '/project')
  857. await ctx.sessionPersistence.create(m)
  858. await ctx.sessionPersistence.append(m.id, oneTurnLog())
  859. const dir = sessionDir(root, m.cwd, m.id)
  860. await writeFile(join(dir, 'metadata.json'), '{}\n')
  861. await writeFile(join(projectDir(root, m.cwd), 'README'), 'project metadata\n')
  862. await mkdir(join(projectDir(root, m.cwd), 'reserved-session'), { recursive: true })
  863. expect(await readdir(dir)).toEqual(expect.arrayContaining(['metadata.json', 'session.jsonl']))
  864. expect((await ctx.sessionPersistence.list()).map(header => header.id)).toContain(m.id)
  865. expect((await ctx.sessionPersistence.load(m.id)).events).toEqual(oneTurnLog())
  866. })
  867. it('rejects the obsolete flat-file layout instead of ignoring stored sessions', async () => {
  868. const m = meta('legacy-flat', '/legacy')
  869. const project = projectDir(root, m.cwd)
  870. const path = join(project, `${encodeSegment(m.id)}.jsonl`)
  871. await mkdir(project, { recursive: true })
  872. await writeFile(path, [
  873. JSON.stringify(toHeaderLine(m)),
  874. ...oneTurnLog().map(event => JSON.stringify(event)),
  875. '',
  876. ].join('\n'))
  877. await expect(ctx.sessionPersistence.load(m.id)).rejects.toThrow(/unsupported flat-file layout/)
  878. await expect(ctx.sessionPersistence.list()).rejects.toThrow(/unsupported flat-file layout/)
  879. })
  880. it('rejects a compressed obsolete flat-file artifact during targeted lookup', async () => {
  881. const m = meta('legacy-compressed-flat', '/legacy')
  882. const project = projectDir(root, m.cwd)
  883. expect(await ctx.sessionPersistence.list()).toEqual([])
  884. await mkdir(project, { recursive: true })
  885. await writeFile(join(project, `${encodeSegment(m.id)}.jsonl.zstd`), 'legacy')
  886. await expect(ctx.sessionPersistence.load(m.id)).rejects.toThrow(/unsupported flat-file layout/)
  887. })
  888. it('list skips empty and non-header session logs (metadata-only read)', async () => {
  889. // A real session…
  890. await ctx.sessionPersistence.create(meta('real', '/p'))
  891. await ctx.sessionPersistence.append(SessionId('real'), oneTurnLog())
  892. // …alongside junk session directories whose fixed transcript is empty or
  893. // lacks a header. Both remain unmaterialized and are skipped.
  894. for (const [id, content] of [
  895. ['empty', ''],
  896. ['notheader', '{"type":"turn/start"}\n'],
  897. ['badjson', 'not json at all\n'],
  898. ] as const) {
  899. const path = rawLogPath(root, undefined, SessionId(id))
  900. await mkdir(sessionDir(root, undefined, SessionId(id)), { recursive: true })
  901. await writeFile(path, content)
  902. }
  903. const ids = (await ctx.sessionPersistence.list()).map(x => x.id).sort()
  904. expect(ids).toEqual(['real'])
  905. })
  906. it('list reads a header line longer than the 8KB read chunk', async () => {
  907. // A tolerated extra field makes this valid header exceed the 8192-byte read buffer, proving
  908. // `readFirstLine` accumulates chunks before `list()` parses it.
  909. const id = SessionId('big')
  910. await mkdir(sessionDir(root, undefined, id), { recursive: true })
  911. const bigHeader = JSON.stringify({ type: 'session', version: 0, id: 'big', createdAt: 1, delegationDepth: 0, pad: 'x'.repeat(9000) })
  912. await writeFile(rawLogPath(root, undefined, id), bigHeader + '\n')
  913. const ids = (await ctx.sessionPersistence.list()).map(x => x.id)
  914. expect(ids).toContain('big')
  915. })
  916. it.each(['sandboxMode', 'approvalPolicy'] as const)('rejects the retired %s header field', (field) => {
  917. const line = { ...toHeaderLine(meta('retired-policy-header')), [field]: 'read-only' }
  918. expect(() => scanLog(Buffer.from(`${JSON.stringify(line)}\n`)))
  919. .toThrow(/retired policy baseline fields/)
  920. })
  921. it('list rejects a header whose cwd does not identify its physical log', async () => {
  922. const m = meta('misplaced', '/stored')
  923. await ctx.sessionPersistence.create(m)
  924. await ctx.sessionPersistence.append(m.id, oneTurnLog())
  925. await rewriteHeader(rawLogPath(root, m.cwd, m.id), (header) => { header.cwd = '/elsewhere' })
  926. await expect(ctx.sessionPersistence.list()).rejects.toThrow(/and cwd identify/)
  927. })
  928. it('accepts an alternate project path only when it identifies the same physical log', async () => {
  929. const m = meta('physical-alias', '/stored')
  930. await ctx.sessionPersistence.create(m)
  931. await ctx.sessionPersistence.append(m.id, oneTurnLog())
  932. const path = rawLogPath(root, m.cwd, m.id)
  933. const aliasCwd = '/alias'
  934. await symlink(
  935. projectDir(root, m.cwd),
  936. projectDir(root, aliasCwd),
  937. process.platform === 'win32' ? 'junction' : 'dir',
  938. )
  939. await rewriteHeader(path, (header) => { header.cwd = aliasCwd })
  940. expect((await ctx.sessionPersistence.load(m.id)).meta.cwd).toBe(aliasCwd)
  941. expect((await ctx.sessionPersistence.list()).map(header => header.id)).toContain(m.id)
  942. })
  943. it('list rejects a session header whose id cannot name a storage path', async () => {
  944. const dir = join(projectDir(root, undefined), 'invalid-id')
  945. await mkdir(dir, { recursive: true })
  946. await writeFile(join(dir, 'session.jsonl'), JSON.stringify({
  947. type: 'session', version: 0, id: '', createdAt: 1, delegationDepth: 0,
  948. }) + '\n')
  949. await expect(ctx.sessionPersistence.list()).rejects.toThrow(/header id cannot name a storage path/)
  950. })
  951. it('load and list reject one id materialized in multiple project directories', async () => {
  952. const id = SessionId('duplicate')
  953. for (const cwd of ['/a', '/b']) {
  954. const m = meta(id, cwd)
  955. await mkdir(sessionDir(root, cwd, id), { recursive: true })
  956. const content = [JSON.stringify(toHeaderLine(m)), ...oneTurnLog().map(event => JSON.stringify(event))].join('\n') + '\n'
  957. await writeFile(rawLogPath(root, cwd, id), content)
  958. }
  959. await expect(ctx.sessionPersistence.load(id)).rejects.toThrow(/appears in multiple project directories/)
  960. await expect(ctx.sessionPersistence.list()).rejects.toThrow(/appears in multiple project directories/)
  961. })
  962. it('a DIFFERENT live session object reusing a disposed id gets its own init (no stale cache)', async () => {
  963. // Session A materializes a log under id "reuse".
  964. const sessFiberA = await ctx.plugin(Object.assign((inner: Context) => {
  965. const a = inner.sessions.create(SessionId('reuse'), { meta: { cwd: '/a' } })
  966. appendLog(a, oneTurnLog())
  967. }, { inject: ['sessions'] }))
  968. // Drain A, then dispose ITS fiber (the live session A is gone) while the
  969. // backend stays loaded.
  970. for (const s of ctx.sessions.list()) await ctx.sessions.flush(s)
  971. await sessFiberA.dispose()
  972. // A new Session object reuses the id. Object-keyed initialization must run independently,
  973. // detect the disk collision, and reject instead of appending through session A's stale cursor.
  974. let b!: Session
  975. await ctx.plugin(Object.assign((inner: Context) => {
  976. b = inner.sessions.create(SessionId('reuse'), { meta: { cwd: '/a' } })
  977. }, { inject: ['sessions'] }))
  978. await expect(ctx.sessions.flush(b)).rejects.toThrow(/already bound to a different live session|already has a persisted log on disk/)
  979. })
  980. it('a no-cwd live session cannot adopt a same-id log from another cwd', async () => {
  981. // Backend 1: materialize a log under id "x" in the cwd "/w" bucket, then
  982. // dispose the WHOLE backend (so backend 2 mounts with an EMPTY states map —
  983. // the HMR/reload path with no tracked collision state).
  984. await ctx.sessionPersistence.create(meta('x', '/w'))
  985. await ctx.sessionPersistence.append(SessionId('x'), oneTurnLog())
  986. await ctx.fiber.dispose()
  987. // Backend 2 creates a no-cwd session whose id exists only in `/w`. The
  988. // stored cwd check rejects instead of grafting no-cwd events onto that log.
  989. const ctx2 = new Context()
  990. await ctx2.plugin(SessionStore)
  991. await ctx2.plugin(SessionPersistenceJsonl, { root, compression: 'none' })
  992. let b!: Session
  993. await ctx2.plugin(Object.assign((inner: Context) => {
  994. b = inner.sessions.create(SessionId('x')) // no cwd
  995. }, { inject: ['sessions'] }))
  996. await expect(ctx2.sessions.flush(b)).rejects.toThrow(/different cwd|id collision/)
  997. // The "/w" log is untouched — no no-cwd events were grafted onto it, and no
  998. // `_no-cwd` log for "x" was created.
  999. const inW = scanLog(await readFile(rawLogPath(root, '/w', SessionId('x'))))
  1000. expect(inW.meta.cwd).toBe('/w')
  1001. expect(inW.events).toHaveLength(6)
  1002. await expect(stat(rawLogPath(root, undefined, SessionId('x')))).rejects.toThrow()
  1003. await ctx2.fiber.dispose()
  1004. })
  1005. it('a seed with matching seq/type/time but DIFFERENT data is rejected (deep prefix compare)', async () => {
  1006. // Materialize and load (ownerless, cursor = 6).
  1007. await ctx.sessionPersistence.create(meta('divergent', '/a'))
  1008. await ctx.sessionPersistence.append(SessionId('divergent'), oneTurnLog())
  1009. await ctx.sessionPersistence.load(SessionId('divergent'))
  1010. // A seed that keeps every seq/type/time but mutates a payload must NOT be
  1011. // accepted as "the same session" — otherwise drain filters those seqs as
  1012. // already persisted and the divergent payload is silently lost.
  1013. const tampered = structuredClone(oneTurnLog())
  1014. const userMsg = tampered[1]
  1015. if (userMsg?.type === 'user/message') {
  1016. (userMsg.data as { content: unknown[] }).content = [{ type: 'text', text: 'DIFFERENT' }]
  1017. }
  1018. let bad!: Session
  1019. await ctx.plugin(Object.assign((inner: Context) => {
  1020. bad = inner.sessions.create(SessionId('divergent'), { seed: tampered, meta: { cwd: '/a' } })
  1021. }, { inject: ['sessions'] }))
  1022. await expect(ctx.sessions.flush(bad)).rejects.toThrow(/do not match this live session|already has a persisted log/)
  1023. })
  1024. it('a second live session reusing a bound id is rejected', async () => {
  1025. // A live session materializes and owns the id.
  1026. const firstFiber = await ctx.plugin(Object.assign((inner: Context) => {
  1027. const a = inner.sessions.create(SessionId('bound'), { meta: { cwd: '/a' } })
  1028. a.append('turn/start', { turn: 1 })
  1029. a.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  1030. }, { inject: ['sessions'] }))
  1031. for (const s of ctx.sessions.list()) await ctx.sessions.flush(s)
  1032. await firstFiber.dispose()
  1033. let second!: Session
  1034. await ctx.plugin(Object.assign((inner: Context) => {
  1035. second = inner.sessions.create(SessionId('bound'), { meta: { cwd: '/a' } })
  1036. }, { inject: ['sessions'] }))
  1037. await expect(ctx.sessions.flush(second))
  1038. .rejects.toThrow(/already bound to a different live session|already has a persisted log|do not match/)
  1039. })
  1040. it('list returns nothing when the root directory does not exist', async () => {
  1041. const ctx2 = new Context()
  1042. await ctx2.plugin(SessionStore)
  1043. await ctx2.plugin(SessionPersistenceJsonl, {
  1044. root: join(root, 'does-not-exist-yet'),
  1045. compression: 'none',
  1046. })
  1047. expect(await ctx2.sessionPersistence.list()).toEqual([])
  1048. await ctx2.fiber.dispose()
  1049. })
  1050. it('plugin load rejects an existing root that is not a directory', async () => {
  1051. const filePath = join(root, 'not-a-dir')
  1052. await writeFile(filePath, 'x')
  1053. const ctx2 = new Context()
  1054. await ctx2.plugin(SessionStore)
  1055. await expect(ctx2.plugin(SessionPersistenceJsonl, { root: filePath, compression: 'none' })).rejects.toThrow(/ENOTDIR/)
  1056. await ctx2.fiber.dispose()
  1057. })
  1058. it('list surfaces a root that becomes unusable after plugin load', async () => {
  1059. await rm(root, { recursive: true })
  1060. await writeFile(root, 'not a directory')
  1061. await expect(ctx.sessionPersistence.list()).rejects.toThrow(/ENOTDIR/)
  1062. })
  1063. it('per-id lookup surfaces non-ENOENT storage errors', async () => {
  1064. const blocker = join(root, 'not-a-directory')
  1065. await writeFile(blocker, 'x')
  1066. const backend = ctx.sessionPersistence as unknown as { exists(path: string): Promise<boolean> }
  1067. await expect(backend.exists(join(blocker, 'child.jsonl'))).rejects.toThrow(/ENOTDIR/)
  1068. })
  1069. it('materialization surfaces a project-directory storage fault', async () => {
  1070. const cwd = '/x'
  1071. const ctx2 = new Context()
  1072. await ctx2.plugin(SessionStore)
  1073. await ctx2.plugin(SessionPersistenceJsonl, { root, compression: 'none' })
  1074. await writeFile(projectDir(root, cwd), 'x') // project path is now a file
  1075. let s!: Session
  1076. await ctx2.plugin(Object.assign((inner: Context) => {
  1077. s = inner.sessions.create(SessionId('exists-fault'), { meta: { cwd } })
  1078. appendClosedTurn(s)
  1079. }, { inject: ['sessions'] }))
  1080. await expect(ctx2.sessions.flush(s)).rejects.toThrow(/EEXIST|ENOTDIR/)
  1081. await ctx2.fiber.dispose()
  1082. })
  1083. it('append() to a disk-only session adopts it and repairs a crash tail', async () => {
  1084. // Persist a session, then corrupt its tail, all through ONE backend.
  1085. const m = meta('disk-append', '/d')
  1086. await ctx.sessionPersistence.create(m)
  1087. await ctx.sessionPersistence.append(m.id, oneTurnLog())
  1088. await writeFile(rawLogPath(root, '/d', m.id), '\n{"partial crash', { flag: 'a' })
  1089. // A FRESH backend with no in-memory state: append directly (no prior load)
  1090. // → append must adopt from disk, and the adopt's load schedules a repair
  1091. // that the same append then performs before writing.
  1092. const ctx2 = new Context()
  1093. await ctx2.plugin(SessionStore)
  1094. await ctx2.plugin(SessionPersistenceJsonl, { root, compression: 'none' })
  1095. await ctx2.sessionPersistence.append(m.id, [
  1096. { type: 'turn/start', seq: 6, time: 9, data: { turn: 2 } },
  1097. { type: 'turn/end', seq: 7, time: 10, data: { turn: 2, reason: { kind: 'completed' } } },
  1098. ] as SessionEvent[])
  1099. const loaded = await ctx2.sessionPersistence.load(m.id)
  1100. expect(loaded.events.map(e => e.seq)).toEqual([0, 1, 2, 3, 4, 5, 6, 7])
  1101. await ctx2.fiber.dispose()
  1102. })
  1103. it('a header-only log (open turn, no turn/end) preserves the open turn on load and closes it', async () => {
  1104. // A session whose only durable content is an unclosed first turn. scanLog
  1105. // preserves the turn/start; loadCore closes it with a synthetic
  1106. // turn/end {interrupted} so the returned log is balanced.
  1107. const m = meta('open-turn', '/h')
  1108. await ctx.sessionPersistence.create(m)
  1109. await ctx.sessionPersistence.append(m.id, [
  1110. { type: 'turn/start', seq: 0, time: 1, data: { turn: 1 } },
  1111. ] as SessionEvent[])
  1112. const { events } = await ctx.sessionPersistence.load(m.id)
  1113. expect(events.map(e => e.type)).toEqual(['turn/start', 'turn/end'])
  1114. const end = events[1]!
  1115. expect(end.type === 'turn/end' && end.data.reason).toEqual({ kind: 'interrupted' })
  1116. })
  1117. it('createCore rejects an id already on disk under a different project directory', async () => {
  1118. // Persist the id under cwd A.
  1119. const a = meta('dup-id', '/projA')
  1120. await ctx.sessionPersistence.create(a)
  1121. await ctx.sessionPersistence.append(a.id, oneTurnLog())
  1122. // A fresh backend creating the SAME id under cwd B must still refuse: load
  1123. // identifies by id across all projects, so a second log would make resume
  1124. // nondeterministic. create scans every project, not just meta.cwd's.
  1125. const ctx2 = new Context()
  1126. await ctx2.plugin(SessionStore)
  1127. await ctx2.plugin(SessionPersistenceJsonl, { root, compression: 'none' })
  1128. await expect(ctx2.sessionPersistence.create(meta('dup-id', '/projB')))
  1129. .rejects.toThrow(/already has a persisted log on disk/)
  1130. await ctx2.fiber.dispose()
  1131. })
  1132. it('flush keeps buffered events when the append fails (no silent loss)', async () => {
  1133. root = await freshRoot()
  1134. const ctx2 = new Context()
  1135. await ctx2.plugin(SessionStore)
  1136. await ctx2.plugin(SessionPersistenceJsonl, { root, compression: 'none' })
  1137. const session = ctx2.sessions.create(SessionId('flush-fail'))
  1138. // A full turn lands in the write-behind buffer.
  1139. session.append('turn/start', { turn: 1 })
  1140. session.append('user/message', createUserMessage({
  1141. content: [{ type: 'text', text: 'hi' }], source: { kind: 'user' },
  1142. }), { surfaceOp: 'append' })
  1143. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  1144. // Make the durable materialize fail on the next flush.
  1145. const backend = ctx2.sessionPersistence as unknown as { materialize: (...args: unknown[]) => Promise<void> }
  1146. const origMat = backend.materialize.bind(backend)
  1147. backend.materialize = () => Promise.reject(new Error('disk full'))
  1148. await expectFlushError(ctx2.sessions.flush(session), /disk full/)
  1149. // The events are STILL buffered (not silently dropped): a retry persists them.
  1150. backend.materialize = origMat
  1151. await ctx2.sessions.flush(session)
  1152. const loaded = await ctx2.sessionPersistence.load(SessionId('flush-fail'))
  1153. expect(loaded.events.map(e => e.seq)).toEqual([0, 1, 2])
  1154. await ctx2.fiber.dispose()
  1155. })
  1156. it('rejects non-JSON event data: BigInt, function, circular, Map, undefined property', async () => {
  1157. const m = meta('serial')
  1158. await ctx.sessionPersistence.create(m)
  1159. const bad = (extra: unknown) => [{
  1160. type: 'user/message',
  1161. seq: 0,
  1162. time: 1,
  1163. data: {
  1164. id: MessageId('invalid-json'),
  1165. role: 'user',
  1166. content: [{ type: 'text', text: 'x' }],
  1167. source: { kind: 'user' },
  1168. extra,
  1169. },
  1170. }] as unknown as SessionEvent[]
  1171. await expect(ctx.sessionPersistence.append(m.id, bad(1n))).rejects.toThrow(/non-JSON-serializable/)
  1172. await expect(ctx.sessionPersistence.append(m.id, bad(() => 0))).rejects.toThrow(/non-JSON-serializable/)
  1173. await expect(ctx.sessionPersistence.append(m.id, bad(Symbol('s')))).rejects.toThrow(/non-JSON-serializable/)
  1174. await expect(ctx.sessionPersistence.append(m.id, bad(new Map()))).rejects.toThrow(/non-JSON-serializable/)
  1175. await expect(ctx.sessionPersistence.append(m.id, bad(undefined))).rejects.toThrow(/non-JSON-serializable/)
  1176. await expect(ctx.sessionPersistence.append(m.id, bad(Infinity))).rejects.toThrow(/non-JSON-serializable/)
  1177. // a circular structure
  1178. const circ: Record<string, unknown> = {}
  1179. circ.self = circ
  1180. await expect(ctx.sessionPersistence.append(m.id, bad(circ))).rejects.toThrow(/non-JSON-serializable/)
  1181. // The session was never materialized by any of the rejected appends.
  1182. expect((await ctx.sessionPersistence.list()).map(h => h.id)).not.toContain(m.id)
  1183. })
  1184. it('accepts well-formed JSON values (null, booleans, nested arrays/objects)', async () => {
  1185. const m = meta('json-ok')
  1186. await ctx.sessionPersistence.create(m)
  1187. const ev = [{ type: 'user/message', seq: 0, time: 1, data: createUserMessage({
  1188. content: [{ type: 'text', text: 'x' }], source: { kind: 'user' }, extra: { a: null, b: true, c: [1, 2, { d: 'nested' }] },
  1189. }) }] as unknown as SessionEvent[]
  1190. await ctx.sessionPersistence.append(m.id, ev)
  1191. expect((await ctx.sessionPersistence.list()).map(h => h.id)).toContain(m.id)
  1192. })
  1193. it('Session.append rejects a non-serializable event at the source (never enters the log)', () => {
  1194. const session = ctx.sessions.create(SessionId('reject-bad'))
  1195. // Serializability is enforced at the source: Session.append throws on a BigInt-bearing
  1196. // event before it enters session.events, so the durable log can never diverge from the live
  1197. // log. The error therefore surfaces synchronously at append, not later during backend flush.
  1198. expect(() => {
  1199. session.append('user/message', { content: [{ type: 'text', text: 'bad' }], source: { kind: 'user' }, bad: 1n } as never, { surfaceOp: 'append' })
  1200. }).toThrow(/non-JSON-serializable/)
  1201. // The bad event was rejected, so the log stayed empty.
  1202. expect(session.events.length).toBe(0)
  1203. })
  1204. })