jsonl.spec.ts 76 KB

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