index.ts 32 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808
  1. /**
  2. * JSONL durable session-persistence backend. It stores a header and contiguous
  3. * events in one append-only file per session, and delegates orchestration to
  4. * {@link PersistenceCoordinator}. Its side-effect-free locator returns the
  5. * absolute per-session log target before materialization.
  6. * @module @deepseek-ai/dsh-session-persistence-jsonl
  7. */
  8. import { Context } from 'cordis'
  9. import z from 'schemastery'
  10. import { readdirSync } from 'node:fs'
  11. import { open, mkdir, readFile, readdir, realpath, link, rm, stat, truncate } from 'node:fs/promises'
  12. import { dirname, join, resolve } from 'node:path'
  13. import { randomBytes } from 'node:crypto'
  14. import {
  15. SessionPersistence, SessionPersistenceRevision, PersistenceCoordinator,
  16. type PersistenceBackend, type SessionLocation, type SessionPersistenceSnapshot,
  17. type StoredPrefix,
  18. } from '@deepseek-ai/dsh-session-persistence'
  19. import type { SessionEvent, SessionId, SessionHeader } from '@deepseek-ai/dsh-session'
  20. import {
  21. encodeSegment, eventLines, logPath, logSuffix, parseHeaderMeta, projectDir, scanLog, sessionDir, toHeaderLine,
  22. type JsonlCompression,
  23. } from './format.ts'
  24. import { compressZstdFrame, decompressZstdFrame, scanZstdFrames } from './zstd.ts'
  25. import { ensureDurableDirectoryWin32, publishNewFileWin32 } from './win32.ts'
  26. export type { JsonlCompression } from './format.ts'
  27. const DEFAULT_COMPRESSION: JsonlCompression = 'zstd'
  28. /** Loader schema for the JSONL artifact's physical encoding. */
  29. export const JsonlCompressionSchema: z<JsonlCompression> = z.union([
  30. z.const('zstd'),
  31. z.const('none'),
  32. ]).default(DEFAULT_COMPRESSION)
  33. /** Plugin config: where the JSONL backend keeps its session logs, and the packed-row write switch. */
  34. export interface Config {
  35. /**
  36. * Root directory for all session files. Required (no default): a default of
  37. * `process.cwd()` would scatter session files as the process's cwd changes
  38. * (bash calls, subprocesses). Sessions group under human-readable project
  39. * directories, then per-session directories. An existing root must be a
  40. * readable directory; an absent root is created on first materialization.
  41. */
  42. root: string
  43. /**
  44. * Write runs of consecutive `assistant/chunk` delta events as packed
  45. * `text-chunks`/`reasoning-chunks`/`tool-call-chunks` rows (lossless,
  46. * ~60% smaller logs measured on a real session). Off by default while
  47. * snapshot fixtures stay in the one-event-per-line layout: recording with
  48. * packing on rewrites every golden `session.jsonl`. READING packed rows is
  49. * unconditional — a log's layout never depends on this switch.
  50. */
  51. packChunks?: boolean
  52. /** Physical encoding; defaults to checksummed Zstandard frames. */
  53. compression?: JsonlCompression
  54. }
  55. /** Opaque coordinator token for replacing bytes recovered from a torn frame. */
  56. interface JsonlTornMarker {
  57. truncateTo: number
  58. recoveredEvents: SessionEvent[]
  59. }
  60. /** Whether a filesystem error means absence; every non-ENOENT failure must surface. */
  61. function isENOENT(error: unknown): boolean {
  62. return (error as NodeJS.ErrnoException | null)?.code === 'ENOENT'
  63. }
  64. /**
  65. * The JSONL persistence backend. Load as a plugin; it registers as
  66. * `ctx.sessionPersistence` and (via the coordinator) installs the write-path
  67. * listeners. Its torn-tail marker carries the byte offset and any events
  68. * recovered from an incomplete final Zstandard frame.
  69. */
  70. export class SessionPersistenceJsonl extends SessionPersistence implements PersistenceBackend<JsonlTornMarker> {
  71. static inject = ['sessions']
  72. static Config: z<Config> = z.object({
  73. root: z.string().required(),
  74. packChunks: z.boolean().default(false),
  75. compression: JsonlCompressionSchema,
  76. })
  77. /**
  78. * Backend label for coordinator diagnostics and effects. It shadows
  79. * `Service.name` without changing the service key captured by the base
  80. * constructor.
  81. */
  82. override readonly name = 'session-persistence-jsonl'
  83. private root: string
  84. private packChunks: boolean
  85. private compression: JsonlCompression
  86. private coordinator: PersistenceCoordinator<JsonlTornMarker>
  87. private rootEncodingCheck: Promise<void> | undefined
  88. constructor(ctx: Context, public config: Config) {
  89. super(ctx)
  90. // Resolve once so later process.cwd() changes cannot split one backend across roots.
  91. this.root = resolve(config.root)
  92. // schemastery (static Config) applied the default before construction;
  93. // the cast records that runtime fact for exactOptionalPropertyTypes.
  94. this.packChunks = (config as Required<Config>).packChunks
  95. this.compression = config.compression ?? DEFAULT_COMPRESSION
  96. this.assertUsableRoot()
  97. this.coordinator = new PersistenceCoordinator<JsonlTornMarker>(this.ctx, this)
  98. }
  99. // Each backend keeps the typed service surface beside its storage hooks;
  100. // extracting these trivial forwards would add an inheritance seam.
  101. /* jscpd:ignore-start */
  102. // --- SessionPersistence service surface (delegated to the coordinator) ---
  103. /** Resolve the absolute target path without touching the filesystem. */
  104. locate(meta: SessionHeader): SessionLocation {
  105. return { kind: 'jsonl', path: logPath(this.root, meta.cwd, meta.id, this.compression) }
  106. }
  107. create(meta: SessionHeader): Promise<void> {
  108. return this.coordinator.create(meta)
  109. }
  110. append(id: SessionId, events: readonly SessionEvent[]): Promise<void> {
  111. return this.coordinator.append(id, events)
  112. }
  113. load(id: SessionId): Promise<{ meta: SessionHeader; events: SessionEvent[] }> {
  114. return this.coordinator.load(id)
  115. }
  116. inspect(id: SessionId, signal?: AbortSignal): Promise<{ meta: SessionHeader; events: SessionEvent[] }> {
  117. return this.coordinator.inspect(id, signal)
  118. }
  119. // One method serves both public `list` and the backend hook; delegating it to
  120. // the coordinator would call this hook recursively.
  121. /* jscpd:ignore-end */
  122. // --- PersistenceBackend hooks (the file-bytes storage primitives) ---
  123. /** Read a stored prefix by id across all project directories when cwd is unknown. */
  124. async loadStored(id: SessionId, signal?: AbortSignal): Promise<StoredPrefix<JsonlTornMarker> | undefined> {
  125. signal?.throwIfAborted()
  126. await this.ensureRootEncoding()
  127. signal?.throwIfAborted()
  128. const path = await this.findLog(id, signal)
  129. if (path === undefined) return undefined
  130. return this.readPrefix(path, id, signal)
  131. }
  132. /**
  133. * Read a stored prefix and convert torn-tail state to the opaque marker the
  134. * coordinator can round-trip without knowing the physical encoding.
  135. */
  136. private async readPrefix(
  137. path: string,
  138. expectedId?: SessionId,
  139. signal?: AbortSignal,
  140. ): Promise<StoredPrefix<JsonlTornMarker>> {
  141. const buffer = await readFile(path, { signal })
  142. signal?.throwIfAborted()
  143. let prefix: StoredPrefix<JsonlTornMarker>
  144. if (this.compression === 'zstd') {
  145. prefix = await this.readZstdPrefix(buffer, signal)
  146. } else {
  147. signal?.throwIfAborted()
  148. const { meta, events, committedBytes } = scanLog(buffer)
  149. signal?.throwIfAborted()
  150. prefix = {
  151. meta,
  152. events,
  153. ...committedBytes < buffer.byteLength
  154. ? { tornMarker: { truncateTo: committedBytes, recoveredEvents: [] } }
  155. : {},
  156. }
  157. }
  158. signal?.throwIfAborted()
  159. await this.assertStoredIdentity(path, prefix.meta, expectedId, signal)
  160. signal?.throwIfAborted()
  161. return prefix
  162. }
  163. /** Decode complete frames and retain complete JSONL records from a torn final frame. */
  164. private async readZstdPrefix(
  165. buffer: Buffer,
  166. signal?: AbortSignal,
  167. ): Promise<StoredPrefix<JsonlTornMarker>> {
  168. signal?.throwIfAborted()
  169. const { frames, tornStart } = scanZstdFrames(buffer)
  170. signal?.throwIfAborted()
  171. if (frames.length === 0) throw new Error('empty or header-less Zstandard session log')
  172. const plaintextFrames: Buffer[] = []
  173. for (const frame of frames) {
  174. let plaintext: Buffer
  175. try {
  176. signal?.throwIfAborted()
  177. plaintext = await decompressZstdFrame(buffer.subarray(frame.start, frame.end))
  178. } catch (error) {
  179. /* v8 ignore next -- decoder failure plus concurrent abort is timing-dependent */
  180. if (signal?.aborted) signal.throwIfAborted()
  181. throw new Error(`corrupt Zstandard session log: frame at byte ${frame.start} failed validation`, { cause: error })
  182. }
  183. signal?.throwIfAborted()
  184. plaintextFrames.push(plaintext)
  185. }
  186. const headerFrame = plaintextFrames[0]
  187. if (headerFrame === undefined || headerFrame.length === 0 || headerFrame.indexOf(0x0A) !== headerFrame.length - 1) {
  188. throw new Error('corrupt Zstandard session log: first frame is not exactly one header line')
  189. }
  190. signal?.throwIfAborted()
  191. const completePlaintext = Buffer.concat(plaintextFrames)
  192. signal?.throwIfAborted()
  193. const completePrefix = scanLog(completePlaintext)
  194. signal?.throwIfAborted()
  195. if (completePrefix.committedBytes !== completePlaintext.length) {
  196. throw new Error('corrupt Zstandard session log: complete frame contains a torn JSONL record')
  197. }
  198. if (tornStart === undefined) {
  199. return { meta: completePrefix.meta, events: completePrefix.events }
  200. }
  201. let recoveredPlaintext: Buffer = Buffer.alloc(0)
  202. try {
  203. signal?.throwIfAborted()
  204. recoveredPlaintext = await decompressZstdFrame(buffer.subarray(tornStart))
  205. } catch {
  206. /* v8 ignore next -- decoder failure plus concurrent abort is timing-dependent */
  207. if (signal?.aborted) signal.throwIfAborted()
  208. // A structurally incomplete final frame may end before Node's decoder can
  209. // emit any plaintext; the complete prior frames remain recoverable.
  210. }
  211. signal?.throwIfAborted()
  212. const recoveredPrefix = scanLog(Buffer.concat([completePlaintext, recoveredPlaintext]))
  213. signal?.throwIfAborted()
  214. /* v8 ignore next 3 -- appending plaintext cannot shorten the already-scanned complete prefix */
  215. if (recoveredPrefix.events.length < completePrefix.events.length) {
  216. throw new Error('corrupt Zstandard session log: recovered prefix does not extend complete frames')
  217. }
  218. return {
  219. meta: recoveredPrefix.meta,
  220. events: recoveredPrefix.events,
  221. tornMarker: {
  222. truncateTo: tornStart,
  223. recoveredEvents: recoveredPrefix.events.slice(completePrefix.events.length),
  224. },
  225. }
  226. }
  227. /** Durably append a batch, lazily materializing the file when not yet present. */
  228. async appendBatch(meta: SessionHeader, events: readonly SessionEvent[], isMaterialized: boolean): Promise<void> {
  229. await this.ensureRootEncoding()
  230. if (isMaterialized) {
  231. await this.appendLines(meta, events)
  232. } else {
  233. await this.materialize(meta, events)
  234. }
  235. }
  236. /**
  237. * Make a crash repair durable: truncate a torn tail, restore complete events
  238. * decoded from it, then append synthetic closers. Two fsync'd steps — the seam
  239. * does not require this to be atomic.
  240. */
  241. async commitRepair(
  242. meta: SessionHeader,
  243. tornMarker: JsonlTornMarker | undefined,
  244. closers: readonly SessionEvent[],
  245. ): Promise<void> {
  246. if (tornMarker !== undefined) await this.repair(meta, tornMarker.truncateTo)
  247. const repairedEvents = [...(tornMarker?.recoveredEvents ?? []), ...closers]
  248. if (repairedEvents.length > 0) await this.appendLines(meta, repairedEvents)
  249. }
  250. /** List valid unique stored sessions' metadata (header line only — no full-log parse). */
  251. async list(signal?: AbortSignal): Promise<SessionHeader[]> {
  252. return (await this.listArtifacts(signal)).map(artifact => artifact.header)
  253. }
  254. /** List metadata plus a stat-derived identity for each append-only log. */
  255. async listSnapshots(signal?: AbortSignal): Promise<SessionPersistenceSnapshot[]> {
  256. const snapshots: SessionPersistenceSnapshot[] = []
  257. for (const artifact of await this.listArtifacts(signal)) {
  258. signal?.throwIfAborted()
  259. try {
  260. const identity = await stat(artifact.path, { bigint: true })
  261. signal?.throwIfAborted()
  262. snapshots.push({
  263. header: artifact.header,
  264. revision: SessionPersistenceRevision([
  265. identity.dev,
  266. identity.ino,
  267. identity.size,
  268. identity.mtimeNs,
  269. identity.ctimeNs,
  270. ].join(':')),
  271. })
  272. } catch (error: unknown) {
  273. signal?.throwIfAborted()
  274. if (!isENOENT(error)) throw error
  275. }
  276. }
  277. signal?.throwIfAborted()
  278. return snapshots
  279. }
  280. private async listArtifacts(signal?: AbortSignal): Promise<Array<{ header: SessionHeader; path: string }>> {
  281. signal?.throwIfAborted()
  282. await this.ensureRootEncoding()
  283. signal?.throwIfAborted()
  284. const artifacts: Array<{ header: SessionHeader; path: string }> = []
  285. const ids = new Set<SessionId>()
  286. for (const project of await this.listProjectDirs(signal)) {
  287. signal?.throwIfAborted()
  288. for (const dir of await this.listSessionDirs(project, signal)) {
  289. signal?.throwIfAborted()
  290. const opposite = join(dir, `session${logSuffix(this.oppositeCompression())}`)
  291. const oppositeExists = await this.exists(opposite)
  292. signal?.throwIfAborted()
  293. if (oppositeExists) throw this.encodingMismatch(opposite)
  294. const path = join(dir, `session${logSuffix(this.compression)}`)
  295. const pathExists = await this.exists(path)
  296. signal?.throwIfAborted()
  297. if (!pathExists) continue
  298. // Read only headers so listing scales with session count, not log size.
  299. const first = this.compression === 'zstd'
  300. ? await this.readFirstZstdLine(path, signal)
  301. : await this.readFirstLine(path, signal)
  302. signal?.throwIfAborted()
  303. if (first === undefined) continue // empty/half-written file
  304. const meta = parseHeaderMeta(first)
  305. if (meta === undefined) continue // not a session header
  306. await this.assertStoredIdentity(path, meta, undefined, signal)
  307. signal?.throwIfAborted()
  308. if (ids.has(meta.id)) {
  309. throw new Error(`duplicate JSONL session id "${meta.id}" appears in multiple project directories`)
  310. }
  311. ids.add(meta.id)
  312. artifacts.push({ header: meta, path })
  313. }
  314. }
  315. signal?.throwIfAborted()
  316. return artifacts
  317. }
  318. // --- materialization / append / repair (file mechanics) ---
  319. /** Atomically write the header line + first batch (temp-write, fsync, publish). */
  320. private async materialize(meta: SessionHeader, events: readonly SessionEvent[]): Promise<void> {
  321. const project = projectDir(this.root, meta.cwd)
  322. const dir = sessionDir(this.root, meta.cwd, meta.id)
  323. const finalPath = logPath(this.root, meta.cwd, meta.id, this.compression)
  324. await this.rejectOppositeArtifact(meta.cwd, meta.id)
  325. const content = await this.encodeMaterialization(meta, events)
  326. /* v8 ignore next -- native Windows coverage exercises this platform dispatch; Linux covers the POSIX peer */
  327. if (process.platform === 'win32') {
  328. await this.materializeWin32(project, dir, finalPath, meta.id, content)
  329. } else {
  330. await this.materializePosix(project, dir, finalPath, meta.id, content)
  331. }
  332. }
  333. /* v8 ignore start -- Windows uses the Win32 durable-publish path; POSIX coverage exercises this peer. */
  334. private async materializePosix(
  335. project: string,
  336. dir: string,
  337. finalPath: string,
  338. id: SessionId,
  339. content: Buffer | string,
  340. ): Promise<void> {
  341. await mkdir(this.root, { recursive: true, mode: 0o700 })
  342. await this.syncDirPosix(dirname(this.root))
  343. await mkdir(project, { recursive: true, mode: 0o700 })
  344. await this.syncDirPosix(this.root)
  345. await mkdir(dir, { recursive: true, mode: 0o700 })
  346. await this.syncDirPosix(project)
  347. await this.rejectExistingLog(finalPath, id)
  348. const tmp = await this.writeSyncedTempFile(finalPath, content)
  349. // Publish via link()+unlink(), NOT rename(): link fails with EEXIST if the
  350. // final path already exists, so two processes materializing the same id
  351. // concurrently cannot clobber each other. rename() would silently overwrite.
  352. let linked = false
  353. try {
  354. await link(tmp, finalPath)
  355. linked = true
  356. } finally {
  357. // Remove an unpublished temp on failure. After publication, defer cleanup
  358. // until the directory entry is durable so cleanup cannot reject a live log.
  359. /* v8 ignore next -- link failure is the TOCTOU/IO race guarded above; not reachable in test */
  360. if (!linked) await rm(tmp, { force: true })
  361. }
  362. // link() succeeded — the log is published. fsync the directory so the new
  363. // entry survives a power loss: the new link is not crash-durable until the
  364. // parent directory's metadata is synced.
  365. await this.syncDirPosix(dir)
  366. // Best-effort temp cleanup: the log is already published and durable, so a
  367. // failure to remove the (now-redundant) temp hard link must NOT reject the
  368. // append. Swallow only the rm failure; nothing else of consequence runs here.
  369. try {
  370. await rm(tmp, { force: true })
  371. } catch {
  372. /* v8 ignore next -- redundant temp link; publish already durable, rm failure is an unreachable IO edge */
  373. }
  374. }
  375. /* v8 ignore stop */
  376. /* v8 ignore start -- native Windows coverage exercises this integration path */
  377. private async materializeWin32(
  378. project: string,
  379. dir: string,
  380. finalPath: string,
  381. id: SessionId,
  382. content: Buffer | string,
  383. ): Promise<void> {
  384. await ensureDurableDirectoryWin32(this.root)
  385. await ensureDurableDirectoryWin32(project)
  386. await ensureDurableDirectoryWin32(dir)
  387. await this.rejectExistingLog(finalPath, id)
  388. const tmp = await this.writeSyncedTempFile(finalPath, content)
  389. try {
  390. await publishNewFileWin32(tmp, finalPath)
  391. } catch (error) {
  392. await rm(tmp, { force: true })
  393. throw error
  394. }
  395. }
  396. /* v8 ignore stop */
  397. private async rejectExistingLog(finalPath: string, id: SessionId): Promise<void> {
  398. // Never publish over an existing committed log: materialize is the first
  399. // write of a session the backend believes is new. A file here means a
  400. // different session shares this id on disk — reject loudly. (createCore
  401. // already guards the create path, so this is unreachable-in-practice TOCTOU
  402. // defense.)
  403. /* v8 ignore next 3 -- createCore guards collisions before materialize; this is a TOCTOU backstop */
  404. if (await this.exists(finalPath)) {
  405. throw new Error(`refusing to materialize "${id}": a log already exists on disk (load/resume it instead)`)
  406. }
  407. }
  408. private async writeSyncedTempFile(finalPath: string, content: Buffer | string): Promise<string> {
  409. const tmp = `${finalPath}.${randomBytes(6).toString('hex')}.tmp`
  410. const handle = await open(tmp, 'wx', 0o600)
  411. try {
  412. await handle.writeFile(content)
  413. await handle.sync()
  414. } finally {
  415. await handle.close()
  416. }
  417. return tmp
  418. }
  419. /** Encode the header and first batch without combining their frame boundaries. */
  420. private async encodeMaterialization(meta: SessionHeader, events: readonly SessionEvent[]): Promise<Buffer | string> {
  421. const header = JSON.stringify(toHeaderLine(meta)) + '\n'
  422. const body = eventLines(events, this.packChunks) + '\n'
  423. if (this.compression === 'none') return header + body
  424. const headerFrame = await compressZstdFrame(header)
  425. const eventFrame = await compressZstdFrame(body)
  426. return Buffer.concat([headerFrame, eventFrame])
  427. }
  428. /** Encode one durable append batch in the configured physical representation. */
  429. private async encodeEventBatch(events: readonly SessionEvent[]): Promise<Buffer | string> {
  430. const body = eventLines(events, this.packChunks) + '\n'
  431. return this.compression === 'zstd' ? compressZstdFrame(body) : body
  432. }
  433. /** fsync a POSIX directory so a just-created/renamed entry is crash-durable. */
  434. /* v8 ignore start -- Windows uses write-through namespace operations; POSIX coverage exercises directory fsync. */
  435. private async syncDirPosix(dir: string): Promise<void> {
  436. const handle = await open(dir, 'r')
  437. try {
  438. await handle.sync()
  439. } finally {
  440. await handle.close()
  441. }
  442. }
  443. /* v8 ignore stop */
  444. /**
  445. * Append and fsync event lines. On a partial write or sync failure, restore the
  446. * previous size before rethrowing because the unchanged cursor will retry the
  447. * batch; leaving partial bytes would create duplicate sequence numbers.
  448. */
  449. private async appendLines(meta: SessionHeader, events: readonly SessionEvent[]): Promise<void> {
  450. const content = await this.encodeEventBatch(events)
  451. const path = logPath(this.root, meta.cwd, meta.id, this.compression)
  452. const handle = await open(path, 'a')
  453. let closed = false
  454. const closeAppendHandle = async (): Promise<void> => {
  455. if (closed) return
  456. closed = true
  457. await handle.close()
  458. }
  459. try {
  460. const { size: before } = await handle.stat()
  461. try {
  462. await handle.writeFile(content)
  463. await handle.sync()
  464. } catch (error) {
  465. try {
  466. await closeAppendHandle()
  467. await this.rollbackAppend(path, before)
  468. } catch (rollbackError) {
  469. throw new AggregateError([error, rollbackError], `failed to roll back append to "${path}"`)
  470. }
  471. throw error
  472. }
  473. } finally {
  474. await closeAppendHandle()
  475. }
  476. }
  477. private async rollbackAppend(path: string, size: number): Promise<void> {
  478. const handle = await open(path, 'r+')
  479. try {
  480. await handle.truncate(size)
  481. await handle.sync()
  482. } finally {
  483. await handle.close()
  484. }
  485. }
  486. /** Truncate the log file to `offset` bytes and fsync (discard the crash tail). */
  487. private async repair(meta: SessionHeader, offset: number): Promise<void> {
  488. const path = logPath(this.root, meta.cwd, meta.id, this.compression)
  489. await truncate(path, offset)
  490. const handle = await open(path, 'r+')
  491. try {
  492. await handle.sync()
  493. } finally {
  494. await handle.close()
  495. }
  496. }
  497. // --- discovery helpers ---
  498. /**
  499. * Read the first newline-terminated line of a file without loading the whole
  500. * file. Returns undefined if the file is empty or has no complete first line.
  501. * Reads in bounded chunks so a huge log costs only the header read.
  502. */
  503. private async readFirstLine(path: string, signal?: AbortSignal): Promise<string | undefined> {
  504. signal?.throwIfAborted()
  505. const handle = await open(path, 'r')
  506. try {
  507. signal?.throwIfAborted()
  508. const chunks: Buffer[] = []
  509. const buf = Buffer.alloc(8192)
  510. for (;;) {
  511. signal?.throwIfAborted()
  512. const { bytesRead } = await handle.read(buf, 0, buf.length, null)
  513. signal?.throwIfAborted()
  514. if (bytesRead === 0) return undefined // EOF with no newline → no complete line
  515. const slice = buf.subarray(0, bytesRead)
  516. const nl = slice.indexOf(0x0a)
  517. if (nl !== -1) {
  518. chunks.push(slice.subarray(0, nl))
  519. signal?.throwIfAborted()
  520. return Buffer.concat(chunks).toString('utf8')
  521. }
  522. chunks.push(Buffer.from(slice))
  523. }
  524. } finally {
  525. await handle.close()
  526. }
  527. }
  528. /** Read and validate only the independently compressed header frame. */
  529. private async readFirstZstdLine(path: string, signal?: AbortSignal): Promise<string | undefined> {
  530. signal?.throwIfAborted()
  531. const handle = await open(path, 'r')
  532. try {
  533. signal?.throwIfAborted()
  534. let content = Buffer.alloc(0)
  535. const chunk = Buffer.alloc(8192)
  536. for (;;) {
  537. signal?.throwIfAborted()
  538. const { bytesRead } = await handle.read(chunk, 0, chunk.length, null)
  539. signal?.throwIfAborted()
  540. if (bytesRead === 0) return undefined
  541. signal?.throwIfAborted()
  542. content = Buffer.concat([content, chunk.subarray(0, bytesRead)])
  543. signal?.throwIfAborted()
  544. const first = scanZstdFrames(content, 1).frames[0]
  545. signal?.throwIfAborted()
  546. if (first === undefined) continue
  547. let plaintext: Buffer
  548. try {
  549. signal?.throwIfAborted()
  550. plaintext = await decompressZstdFrame(content.subarray(first.start, first.end))
  551. } catch (error) {
  552. /* v8 ignore next -- decoder failure plus concurrent abort is timing-dependent */
  553. if (signal?.aborted) signal.throwIfAborted()
  554. throw new Error('corrupt Zstandard session log: header frame failed validation', { cause: error })
  555. }
  556. signal?.throwIfAborted()
  557. if (plaintext.length === 0 || plaintext.indexOf(0x0A) !== plaintext.length - 1) {
  558. throw new Error('corrupt Zstandard session log: first frame is not exactly one header line')
  559. }
  560. return plaintext.subarray(0, -1).toString('utf8')
  561. }
  562. } finally {
  563. await handle.close()
  564. }
  565. }
  566. /** Find the unique physical log for an id across every project directory. */
  567. private async findLog(id: SessionId, signal?: AbortSignal): Promise<string | undefined> {
  568. const matches: string[] = []
  569. for (const project of await this.listProjectDirs(signal)) {
  570. signal?.throwIfAborted()
  571. await this.rejectLegacyFlatArtifact(project, id, signal)
  572. signal?.throwIfAborted()
  573. const dir = join(project, encodeSegment(id))
  574. const path = join(dir, `session${logSuffix(this.compression)}`)
  575. const opposite = join(dir, `session${logSuffix(this.oppositeCompression())}`)
  576. const oppositeExists = await this.exists(opposite)
  577. signal?.throwIfAborted()
  578. if (oppositeExists) throw this.encodingMismatch(opposite)
  579. const pathExists = await this.exists(path)
  580. signal?.throwIfAborted()
  581. if (pathExists) matches.push(path)
  582. }
  583. if (matches.length > 1) {
  584. throw new Error(`duplicate JSONL session id "${id}" appears in multiple project directories`)
  585. }
  586. signal?.throwIfAborted()
  587. return matches[0]
  588. }
  589. /** Require an existing configured root to be a readable directory. */
  590. private assertUsableRoot(): void {
  591. try {
  592. readdirSync(this.root)
  593. } catch (error) {
  594. if (isENOENT(error)) return
  595. throw error
  596. }
  597. }
  598. /** Reject metadata that does not identify the selected physical log. */
  599. private async assertStoredIdentity(
  600. path: string,
  601. meta: SessionHeader,
  602. expectedId?: SessionId,
  603. signal?: AbortSignal,
  604. ): Promise<void> {
  605. signal?.throwIfAborted()
  606. if (expectedId !== undefined && meta.id !== expectedId) {
  607. throw new Error(`corrupt session log "${path}": requested id "${expectedId}" does not match header id "${meta.id}"`)
  608. }
  609. let expectedPath: string
  610. try {
  611. expectedPath = logPath(this.root, meta.cwd, meta.id, this.compression)
  612. } catch (error) {
  613. throw new Error(`corrupt session log "${path}": header id cannot name a storage path`, { cause: error })
  614. }
  615. if (path !== expectedPath && !await this.sameFile(path, expectedPath, signal)) {
  616. throw new Error(`corrupt session log "${path}": header id "${meta.id}" and cwd identify "${expectedPath}"`)
  617. }
  618. signal?.throwIfAborted()
  619. }
  620. /**
  621. * Whether two path spellings resolve to the same physical file. This admits
  622. * case aliases on case-insensitive filesystems without weakening identity
  623. * checks on case-sensitive stores.
  624. */
  625. private async sameFile(path: string, expectedPath: string, signal?: AbortSignal): Promise<boolean> {
  626. signal?.throwIfAborted()
  627. try {
  628. const [actual, expected] = await Promise.all([realpath(path), realpath(expectedPath)])
  629. signal?.throwIfAborted()
  630. return actual === expected
  631. } catch (error) {
  632. signal?.throwIfAborted()
  633. /* v8 ignore else -- non-ENOENT realpath failures require an external permission or I/O fault */
  634. if (isENOENT(error)) return false
  635. /* v8 ignore next -- non-ENOENT realpath failures are external I/O faults, propagated unchanged */
  636. throw error
  637. }
  638. }
  639. /** The human-readable project directories under the configured root. */
  640. private async listProjectDirs(signal?: AbortSignal): Promise<string[]> {
  641. try {
  642. signal?.throwIfAborted()
  643. const entries = await readdir(this.root, { withFileTypes: true })
  644. signal?.throwIfAborted()
  645. return entries.filter(e => e.isDirectory()).map(e => join(this.root, e.name))
  646. } catch (error) {
  647. // Only an absent root means no sessions; rethrow every other I/O failure.
  648. if (isENOENT(error)) return []
  649. throw error
  650. }
  651. }
  652. /** List session-owned directories and reject the obsolete flat-file layout. */
  653. private async listSessionDirs(project: string, signal?: AbortSignal): Promise<string[]> {
  654. signal?.throwIfAborted()
  655. const entries = await readdir(project, { withFileTypes: true })
  656. signal?.throwIfAborted()
  657. const legacy = entries.find(entry =>
  658. entry.isFile() && (entry.name.endsWith('.jsonl') || entry.name.endsWith('.jsonl.zstd')))
  659. if (legacy !== undefined) throw this.legacyLayout(join(project, legacy.name))
  660. return entries.filter(entry => entry.isDirectory()).map(entry => join(project, entry.name))
  661. }
  662. /** Reject a root that already belongs to the other physical encoding. */
  663. private ensureRootEncoding(): Promise<void> {
  664. this.rootEncodingCheck ??= this.checkRootEncoding()
  665. return this.rootEncodingCheck
  666. }
  667. private async checkRootEncoding(): Promise<void> {
  668. for (const project of await this.listProjectDirs()) {
  669. for (const dir of await this.listSessionDirs(project)) {
  670. const incompatible = join(dir, `session${logSuffix(this.oppositeCompression())}`)
  671. if (await this.exists(incompatible)) throw this.encodingMismatch(incompatible)
  672. }
  673. }
  674. }
  675. private async rejectLegacyFlatArtifact(
  676. project: string,
  677. id: SessionId,
  678. signal?: AbortSignal,
  679. ): Promise<void> {
  680. signal?.throwIfAborted()
  681. const encoded = encodeSegment(id)
  682. for (const compression of ['zstd', 'none'] as const) {
  683. const path = join(project, encoded + logSuffix(compression))
  684. const artifactExists = await this.exists(path)
  685. signal?.throwIfAborted()
  686. if (artifactExists) throw this.legacyLayout(path)
  687. }
  688. }
  689. private async rejectOppositeArtifact(cwd: string | undefined, id: SessionId): Promise<void> {
  690. const path = logPath(this.root, cwd, id, this.oppositeCompression())
  691. if (await this.exists(path)) throw this.encodingMismatch(path)
  692. }
  693. private oppositeCompression(): JsonlCompression {
  694. return this.compression === 'zstd' ? 'none' : 'zstd'
  695. }
  696. private encodingMismatch(path: string): Error {
  697. return new Error(
  698. `session artifact ${JSON.stringify(path)} uses ${logSuffix(this.oppositeCompression())}, `
  699. + `but this backend is configured for compression ${JSON.stringify(this.compression)}; `
  700. + 'use a separate root or select the matching compression mode',
  701. )
  702. }
  703. private legacyLayout(path: string): Error {
  704. return new Error(
  705. `session artifact ${JSON.stringify(path)} uses the unsupported flat-file layout; `
  706. + 'use a separate root or move it into a project/session directory before loading',
  707. )
  708. }
  709. private async exists(path: string): Promise<boolean> {
  710. try {
  711. const handle = await open(path, 'r')
  712. await handle.close()
  713. return true
  714. } catch (error) {
  715. // Only ENOENT means absent. A permission/I/O error must surface rather
  716. // than letting load or collision checks proceed under false absence.
  717. // Windows reports ENOENT, not ENOTDIR, for `regular-file/child`; verify
  718. // the immediate parent so a blocked session directory remains a storage fault.
  719. /* v8 ignore else -- Windows reports file-valued parents as ENOENT; POSIX covers direct ENOTDIR. */
  720. if (isENOENT(error)) {
  721. await this.assertLogParentAllowsAbsence(path)
  722. return false
  723. }
  724. /* v8 ignore next -- Windows repairs ENOTDIR from ENOENT above; POSIX covers direct ENOTDIR. */
  725. throw error
  726. }
  727. }
  728. /* v8 ignore start -- native Windows coverage exercises this repair; POSIX open reports ENOTDIR before this point. */
  729. private async assertLogParentAllowsAbsence(path: string): Promise<void> {
  730. try {
  731. const parent = dirname(path)
  732. const info = await stat(parent)
  733. if (info.isDirectory()) return
  734. const error = new Error(`ENOTDIR: parent path exists but is not a directory: ${parent}`) as NodeJS.ErrnoException
  735. error.code = 'ENOTDIR'
  736. error.path = parent
  737. throw error
  738. } catch (error) {
  739. if (isENOENT(error)) return
  740. throw error
  741. }
  742. }
  743. /* v8 ignore stop */
  744. }
  745. export default SessionPersistenceJsonl