index.ts 38 KB

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