jsonl.spec.ts 56 KB

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