| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808 |
- /**
- * JSONL durable session-persistence backend. It stores a header and contiguous
- * events in one append-only file per session, and delegates orchestration to
- * {@link PersistenceCoordinator}. Its side-effect-free locator returns the
- * absolute per-session log target before materialization.
- * @module @deepseek-ai/dsh-session-persistence-jsonl
- */
- import { Context } from 'cordis'
- import z from 'schemastery'
- import { readdirSync } from 'node:fs'
- import { open, mkdir, readFile, readdir, realpath, link, rm, stat, truncate } from 'node:fs/promises'
- import { dirname, join, resolve } from 'node:path'
- import { randomBytes } from 'node:crypto'
- import {
- SessionPersistence, SessionPersistenceRevision, PersistenceCoordinator,
- type PersistenceBackend, type SessionLocation, type SessionPersistenceSnapshot,
- type StoredPrefix,
- } from '@deepseek-ai/dsh-session-persistence'
- import type { SessionEvent, SessionId, SessionHeader } from '@deepseek-ai/dsh-session'
- import {
- encodeSegment, eventLines, logPath, logSuffix, parseHeaderMeta, projectDir, scanLog, sessionDir, toHeaderLine,
- type JsonlCompression,
- } from './format.ts'
- import { compressZstdFrame, decompressZstdFrame, scanZstdFrames } from './zstd.ts'
- import { ensureDurableDirectoryWin32, publishNewFileWin32 } from './win32.ts'
- export type { JsonlCompression } from './format.ts'
- const DEFAULT_COMPRESSION: JsonlCompression = 'zstd'
- /** Loader schema for the JSONL artifact's physical encoding. */
- export const JsonlCompressionSchema: z<JsonlCompression> = z.union([
- z.const('zstd'),
- z.const('none'),
- ]).default(DEFAULT_COMPRESSION)
- /** Plugin config: where the JSONL backend keeps its session logs, and the packed-row write switch. */
- export interface Config {
- /**
- * Root directory for all session files. Required (no default): a default of
- * `process.cwd()` would scatter session files as the process's cwd changes
- * (bash calls, subprocesses). Sessions group under human-readable project
- * directories, then per-session directories. An existing root must be a
- * readable directory; an absent root is created on first materialization.
- */
- root: string
- /**
- * Write runs of consecutive `assistant/chunk` delta events as packed
- * `text-chunks`/`reasoning-chunks`/`tool-call-chunks` rows (lossless,
- * ~60% smaller logs measured on a real session). Off by default while
- * snapshot fixtures stay in the one-event-per-line layout: recording with
- * packing on rewrites every golden `session.jsonl`. READING packed rows is
- * unconditional — a log's layout never depends on this switch.
- */
- packChunks?: boolean
- /** Physical encoding; defaults to checksummed Zstandard frames. */
- compression?: JsonlCompression
- }
- /** Opaque coordinator token for replacing bytes recovered from a torn frame. */
- interface JsonlTornMarker {
- truncateTo: number
- recoveredEvents: SessionEvent[]
- }
- /** Whether a filesystem error means absence; every non-ENOENT failure must surface. */
- function isENOENT(error: unknown): boolean {
- return (error as NodeJS.ErrnoException | null)?.code === 'ENOENT'
- }
- /**
- * The JSONL persistence backend. Load as a plugin; it registers as
- * `ctx.sessionPersistence` and (via the coordinator) installs the write-path
- * listeners. Its torn-tail marker carries the byte offset and any events
- * recovered from an incomplete final Zstandard frame.
- */
- export class SessionPersistenceJsonl extends SessionPersistence implements PersistenceBackend<JsonlTornMarker> {
- static inject = ['sessions']
- static Config: z<Config> = z.object({
- root: z.string().required(),
- packChunks: z.boolean().default(false),
- compression: JsonlCompressionSchema,
- })
- /**
- * Backend label for coordinator diagnostics and effects. It shadows
- * `Service.name` without changing the service key captured by the base
- * constructor.
- */
- override readonly name = 'session-persistence-jsonl'
- private root: string
- private packChunks: boolean
- private compression: JsonlCompression
- private coordinator: PersistenceCoordinator<JsonlTornMarker>
- private rootEncodingCheck: Promise<void> | undefined
- constructor(ctx: Context, public config: Config) {
- super(ctx)
- // Resolve once so later process.cwd() changes cannot split one backend across roots.
- this.root = resolve(config.root)
- // schemastery (static Config) applied the default before construction;
- // the cast records that runtime fact for exactOptionalPropertyTypes.
- this.packChunks = (config as Required<Config>).packChunks
- this.compression = config.compression ?? DEFAULT_COMPRESSION
- this.assertUsableRoot()
- this.coordinator = new PersistenceCoordinator<JsonlTornMarker>(this.ctx, this)
- }
- // Each backend keeps the typed service surface beside its storage hooks;
- // extracting these trivial forwards would add an inheritance seam.
- /* jscpd:ignore-start */
- // --- SessionPersistence service surface (delegated to the coordinator) ---
- /** Resolve the absolute target path without touching the filesystem. */
- locate(meta: SessionHeader): SessionLocation {
- return { kind: 'jsonl', path: logPath(this.root, meta.cwd, meta.id, this.compression) }
- }
- create(meta: SessionHeader): Promise<void> {
- return this.coordinator.create(meta)
- }
- append(id: SessionId, events: readonly SessionEvent[]): Promise<void> {
- return this.coordinator.append(id, events)
- }
- load(id: SessionId): Promise<{ meta: SessionHeader; events: SessionEvent[] }> {
- return this.coordinator.load(id)
- }
- inspect(id: SessionId, signal?: AbortSignal): Promise<{ meta: SessionHeader; events: SessionEvent[] }> {
- return this.coordinator.inspect(id, signal)
- }
- // One method serves both public `list` and the backend hook; delegating it to
- // the coordinator would call this hook recursively.
- /* jscpd:ignore-end */
- // --- PersistenceBackend hooks (the file-bytes storage primitives) ---
- /** Read a stored prefix by id across all project directories when cwd is unknown. */
- async loadStored(id: SessionId, signal?: AbortSignal): Promise<StoredPrefix<JsonlTornMarker> | undefined> {
- signal?.throwIfAborted()
- await this.ensureRootEncoding()
- signal?.throwIfAborted()
- const path = await this.findLog(id, signal)
- if (path === undefined) return undefined
- return this.readPrefix(path, id, signal)
- }
- /**
- * Read a stored prefix and convert torn-tail state to the opaque marker the
- * coordinator can round-trip without knowing the physical encoding.
- */
- private async readPrefix(
- path: string,
- expectedId?: SessionId,
- signal?: AbortSignal,
- ): Promise<StoredPrefix<JsonlTornMarker>> {
- const buffer = await readFile(path, { signal })
- signal?.throwIfAborted()
- let prefix: StoredPrefix<JsonlTornMarker>
- if (this.compression === 'zstd') {
- prefix = await this.readZstdPrefix(buffer, signal)
- } else {
- signal?.throwIfAborted()
- const { meta, events, committedBytes } = scanLog(buffer)
- signal?.throwIfAborted()
- prefix = {
- meta,
- events,
- ...committedBytes < buffer.byteLength
- ? { tornMarker: { truncateTo: committedBytes, recoveredEvents: [] } }
- : {},
- }
- }
- signal?.throwIfAborted()
- await this.assertStoredIdentity(path, prefix.meta, expectedId, signal)
- signal?.throwIfAborted()
- return prefix
- }
- /** Decode complete frames and retain complete JSONL records from a torn final frame. */
- private async readZstdPrefix(
- buffer: Buffer,
- signal?: AbortSignal,
- ): Promise<StoredPrefix<JsonlTornMarker>> {
- signal?.throwIfAborted()
- const { frames, tornStart } = scanZstdFrames(buffer)
- signal?.throwIfAborted()
- if (frames.length === 0) throw new Error('empty or header-less Zstandard session log')
- const plaintextFrames: Buffer[] = []
- for (const frame of frames) {
- let plaintext: Buffer
- try {
- signal?.throwIfAborted()
- plaintext = await decompressZstdFrame(buffer.subarray(frame.start, frame.end))
- } catch (error) {
- /* v8 ignore next -- decoder failure plus concurrent abort is timing-dependent */
- if (signal?.aborted) signal.throwIfAborted()
- throw new Error(`corrupt Zstandard session log: frame at byte ${frame.start} failed validation`, { cause: error })
- }
- signal?.throwIfAborted()
- plaintextFrames.push(plaintext)
- }
- const headerFrame = plaintextFrames[0]
- if (headerFrame === undefined || headerFrame.length === 0 || headerFrame.indexOf(0x0A) !== headerFrame.length - 1) {
- throw new Error('corrupt Zstandard session log: first frame is not exactly one header line')
- }
- signal?.throwIfAborted()
- const completePlaintext = Buffer.concat(plaintextFrames)
- signal?.throwIfAborted()
- const completePrefix = scanLog(completePlaintext)
- signal?.throwIfAborted()
- if (completePrefix.committedBytes !== completePlaintext.length) {
- throw new Error('corrupt Zstandard session log: complete frame contains a torn JSONL record')
- }
- if (tornStart === undefined) {
- return { meta: completePrefix.meta, events: completePrefix.events }
- }
- let recoveredPlaintext: Buffer = Buffer.alloc(0)
- try {
- signal?.throwIfAborted()
- recoveredPlaintext = await decompressZstdFrame(buffer.subarray(tornStart))
- } catch {
- /* v8 ignore next -- decoder failure plus concurrent abort is timing-dependent */
- if (signal?.aborted) signal.throwIfAborted()
- // A structurally incomplete final frame may end before Node's decoder can
- // emit any plaintext; the complete prior frames remain recoverable.
- }
- signal?.throwIfAborted()
- const recoveredPrefix = scanLog(Buffer.concat([completePlaintext, recoveredPlaintext]))
- signal?.throwIfAborted()
- /* v8 ignore next 3 -- appending plaintext cannot shorten the already-scanned complete prefix */
- if (recoveredPrefix.events.length < completePrefix.events.length) {
- throw new Error('corrupt Zstandard session log: recovered prefix does not extend complete frames')
- }
- return {
- meta: recoveredPrefix.meta,
- events: recoveredPrefix.events,
- tornMarker: {
- truncateTo: tornStart,
- recoveredEvents: recoveredPrefix.events.slice(completePrefix.events.length),
- },
- }
- }
- /** Durably append a batch, lazily materializing the file when not yet present. */
- async appendBatch(meta: SessionHeader, events: readonly SessionEvent[], isMaterialized: boolean): Promise<void> {
- await this.ensureRootEncoding()
- if (isMaterialized) {
- await this.appendLines(meta, events)
- } else {
- await this.materialize(meta, events)
- }
- }
- /**
- * Make a crash repair durable: truncate a torn tail, restore complete events
- * decoded from it, then append synthetic closers. Two fsync'd steps — the seam
- * does not require this to be atomic.
- */
- async commitRepair(
- meta: SessionHeader,
- tornMarker: JsonlTornMarker | undefined,
- closers: readonly SessionEvent[],
- ): Promise<void> {
- if (tornMarker !== undefined) await this.repair(meta, tornMarker.truncateTo)
- const repairedEvents = [...(tornMarker?.recoveredEvents ?? []), ...closers]
- if (repairedEvents.length > 0) await this.appendLines(meta, repairedEvents)
- }
- /** List valid unique stored sessions' metadata (header line only — no full-log parse). */
- async list(signal?: AbortSignal): Promise<SessionHeader[]> {
- return (await this.listArtifacts(signal)).map(artifact => artifact.header)
- }
- /** List metadata plus a stat-derived identity for each append-only log. */
- async listSnapshots(signal?: AbortSignal): Promise<SessionPersistenceSnapshot[]> {
- const snapshots: SessionPersistenceSnapshot[] = []
- for (const artifact of await this.listArtifacts(signal)) {
- signal?.throwIfAborted()
- try {
- const identity = await stat(artifact.path, { bigint: true })
- signal?.throwIfAborted()
- snapshots.push({
- header: artifact.header,
- revision: SessionPersistenceRevision([
- identity.dev,
- identity.ino,
- identity.size,
- identity.mtimeNs,
- identity.ctimeNs,
- ].join(':')),
- })
- } catch (error: unknown) {
- signal?.throwIfAborted()
- if (!isENOENT(error)) throw error
- }
- }
- signal?.throwIfAborted()
- return snapshots
- }
- private async listArtifacts(signal?: AbortSignal): Promise<Array<{ header: SessionHeader; path: string }>> {
- signal?.throwIfAborted()
- await this.ensureRootEncoding()
- signal?.throwIfAborted()
- const artifacts: Array<{ header: SessionHeader; path: string }> = []
- const ids = new Set<SessionId>()
- for (const project of await this.listProjectDirs(signal)) {
- signal?.throwIfAborted()
- for (const dir of await this.listSessionDirs(project, signal)) {
- signal?.throwIfAborted()
- const opposite = join(dir, `session${logSuffix(this.oppositeCompression())}`)
- const oppositeExists = await this.exists(opposite)
- signal?.throwIfAborted()
- if (oppositeExists) throw this.encodingMismatch(opposite)
- const path = join(dir, `session${logSuffix(this.compression)}`)
- const pathExists = await this.exists(path)
- signal?.throwIfAborted()
- if (!pathExists) continue
- // Read only headers so listing scales with session count, not log size.
- const first = this.compression === 'zstd'
- ? await this.readFirstZstdLine(path, signal)
- : await this.readFirstLine(path, signal)
- signal?.throwIfAborted()
- if (first === undefined) continue // empty/half-written file
- const meta = parseHeaderMeta(first)
- if (meta === undefined) continue // not a session header
- await this.assertStoredIdentity(path, meta, undefined, signal)
- signal?.throwIfAborted()
- if (ids.has(meta.id)) {
- throw new Error(`duplicate JSONL session id "${meta.id}" appears in multiple project directories`)
- }
- ids.add(meta.id)
- artifacts.push({ header: meta, path })
- }
- }
- signal?.throwIfAborted()
- return artifacts
- }
- // --- materialization / append / repair (file mechanics) ---
- /** Atomically write the header line + first batch (temp-write, fsync, publish). */
- private async materialize(meta: SessionHeader, events: readonly SessionEvent[]): Promise<void> {
- const project = projectDir(this.root, meta.cwd)
- const dir = sessionDir(this.root, meta.cwd, meta.id)
- const finalPath = logPath(this.root, meta.cwd, meta.id, this.compression)
- await this.rejectOppositeArtifact(meta.cwd, meta.id)
- const content = await this.encodeMaterialization(meta, events)
- /* v8 ignore next -- native Windows coverage exercises this platform dispatch; Linux covers the POSIX peer */
- if (process.platform === 'win32') {
- await this.materializeWin32(project, dir, finalPath, meta.id, content)
- } else {
- await this.materializePosix(project, dir, finalPath, meta.id, content)
- }
- }
- /* v8 ignore start -- Windows uses the Win32 durable-publish path; POSIX coverage exercises this peer. */
- private async materializePosix(
- project: string,
- dir: string,
- finalPath: string,
- id: SessionId,
- content: Buffer | string,
- ): Promise<void> {
- await mkdir(this.root, { recursive: true, mode: 0o700 })
- await this.syncDirPosix(dirname(this.root))
- await mkdir(project, { recursive: true, mode: 0o700 })
- await this.syncDirPosix(this.root)
- await mkdir(dir, { recursive: true, mode: 0o700 })
- await this.syncDirPosix(project)
- await this.rejectExistingLog(finalPath, id)
- const tmp = await this.writeSyncedTempFile(finalPath, content)
- // Publish via link()+unlink(), NOT rename(): link fails with EEXIST if the
- // final path already exists, so two processes materializing the same id
- // concurrently cannot clobber each other. rename() would silently overwrite.
- let linked = false
- try {
- await link(tmp, finalPath)
- linked = true
- } finally {
- // Remove an unpublished temp on failure. After publication, defer cleanup
- // until the directory entry is durable so cleanup cannot reject a live log.
- /* v8 ignore next -- link failure is the TOCTOU/IO race guarded above; not reachable in test */
- if (!linked) await rm(tmp, { force: true })
- }
- // link() succeeded — the log is published. fsync the directory so the new
- // entry survives a power loss: the new link is not crash-durable until the
- // parent directory's metadata is synced.
- await this.syncDirPosix(dir)
- // Best-effort temp cleanup: the log is already published and durable, so a
- // failure to remove the (now-redundant) temp hard link must NOT reject the
- // append. Swallow only the rm failure; nothing else of consequence runs here.
- try {
- await rm(tmp, { force: true })
- } catch {
- /* v8 ignore next -- redundant temp link; publish already durable, rm failure is an unreachable IO edge */
- }
- }
- /* v8 ignore stop */
- /* v8 ignore start -- native Windows coverage exercises this integration path */
- private async materializeWin32(
- project: string,
- dir: string,
- finalPath: string,
- id: SessionId,
- content: Buffer | string,
- ): Promise<void> {
- await ensureDurableDirectoryWin32(this.root)
- await ensureDurableDirectoryWin32(project)
- await ensureDurableDirectoryWin32(dir)
- await this.rejectExistingLog(finalPath, id)
- const tmp = await this.writeSyncedTempFile(finalPath, content)
- try {
- await publishNewFileWin32(tmp, finalPath)
- } catch (error) {
- await rm(tmp, { force: true })
- throw error
- }
- }
- /* v8 ignore stop */
- private async rejectExistingLog(finalPath: string, id: SessionId): Promise<void> {
- // Never publish over an existing committed log: materialize is the first
- // write of a session the backend believes is new. A file here means a
- // different session shares this id on disk — reject loudly. (createCore
- // already guards the create path, so this is unreachable-in-practice TOCTOU
- // defense.)
- /* v8 ignore next 3 -- createCore guards collisions before materialize; this is a TOCTOU backstop */
- if (await this.exists(finalPath)) {
- throw new Error(`refusing to materialize "${id}": a log already exists on disk (load/resume it instead)`)
- }
- }
- private async writeSyncedTempFile(finalPath: string, content: Buffer | string): Promise<string> {
- const tmp = `${finalPath}.${randomBytes(6).toString('hex')}.tmp`
- const handle = await open(tmp, 'wx', 0o600)
- try {
- await handle.writeFile(content)
- await handle.sync()
- } finally {
- await handle.close()
- }
- return tmp
- }
- /** Encode the header and first batch without combining their frame boundaries. */
- private async encodeMaterialization(meta: SessionHeader, events: readonly SessionEvent[]): Promise<Buffer | string> {
- const header = JSON.stringify(toHeaderLine(meta)) + '\n'
- const body = eventLines(events, this.packChunks) + '\n'
- if (this.compression === 'none') return header + body
- const headerFrame = await compressZstdFrame(header)
- const eventFrame = await compressZstdFrame(body)
- return Buffer.concat([headerFrame, eventFrame])
- }
- /** Encode one durable append batch in the configured physical representation. */
- private async encodeEventBatch(events: readonly SessionEvent[]): Promise<Buffer | string> {
- const body = eventLines(events, this.packChunks) + '\n'
- return this.compression === 'zstd' ? compressZstdFrame(body) : body
- }
- /** fsync a POSIX directory so a just-created/renamed entry is crash-durable. */
- /* v8 ignore start -- Windows uses write-through namespace operations; POSIX coverage exercises directory fsync. */
- private async syncDirPosix(dir: string): Promise<void> {
- const handle = await open(dir, 'r')
- try {
- await handle.sync()
- } finally {
- await handle.close()
- }
- }
- /* v8 ignore stop */
- /**
- * Append and fsync event lines. On a partial write or sync failure, restore the
- * previous size before rethrowing because the unchanged cursor will retry the
- * batch; leaving partial bytes would create duplicate sequence numbers.
- */
- private async appendLines(meta: SessionHeader, events: readonly SessionEvent[]): Promise<void> {
- const content = await this.encodeEventBatch(events)
- const path = logPath(this.root, meta.cwd, meta.id, this.compression)
- const handle = await open(path, 'a')
- let closed = false
- const closeAppendHandle = async (): Promise<void> => {
- if (closed) return
- closed = true
- await handle.close()
- }
- try {
- const { size: before } = await handle.stat()
- try {
- await handle.writeFile(content)
- await handle.sync()
- } catch (error) {
- try {
- await closeAppendHandle()
- await this.rollbackAppend(path, before)
- } catch (rollbackError) {
- throw new AggregateError([error, rollbackError], `failed to roll back append to "${path}"`)
- }
- throw error
- }
- } finally {
- await closeAppendHandle()
- }
- }
- private async rollbackAppend(path: string, size: number): Promise<void> {
- const handle = await open(path, 'r+')
- try {
- await handle.truncate(size)
- await handle.sync()
- } finally {
- await handle.close()
- }
- }
- /** Truncate the log file to `offset` bytes and fsync (discard the crash tail). */
- private async repair(meta: SessionHeader, offset: number): Promise<void> {
- const path = logPath(this.root, meta.cwd, meta.id, this.compression)
- await truncate(path, offset)
- const handle = await open(path, 'r+')
- try {
- await handle.sync()
- } finally {
- await handle.close()
- }
- }
- // --- discovery helpers ---
- /**
- * Read the first newline-terminated line of a file without loading the whole
- * file. Returns undefined if the file is empty or has no complete first line.
- * Reads in bounded chunks so a huge log costs only the header read.
- */
- private async readFirstLine(path: string, signal?: AbortSignal): Promise<string | undefined> {
- signal?.throwIfAborted()
- const handle = await open(path, 'r')
- try {
- signal?.throwIfAborted()
- const chunks: Buffer[] = []
- const buf = Buffer.alloc(8192)
- for (;;) {
- signal?.throwIfAborted()
- const { bytesRead } = await handle.read(buf, 0, buf.length, null)
- signal?.throwIfAborted()
- if (bytesRead === 0) return undefined // EOF with no newline → no complete line
- const slice = buf.subarray(0, bytesRead)
- const nl = slice.indexOf(0x0a)
- if (nl !== -1) {
- chunks.push(slice.subarray(0, nl))
- signal?.throwIfAborted()
- return Buffer.concat(chunks).toString('utf8')
- }
- chunks.push(Buffer.from(slice))
- }
- } finally {
- await handle.close()
- }
- }
- /** Read and validate only the independently compressed header frame. */
- private async readFirstZstdLine(path: string, signal?: AbortSignal): Promise<string | undefined> {
- signal?.throwIfAborted()
- const handle = await open(path, 'r')
- try {
- signal?.throwIfAborted()
- let content = Buffer.alloc(0)
- const chunk = Buffer.alloc(8192)
- for (;;) {
- signal?.throwIfAborted()
- const { bytesRead } = await handle.read(chunk, 0, chunk.length, null)
- signal?.throwIfAborted()
- if (bytesRead === 0) return undefined
- signal?.throwIfAborted()
- content = Buffer.concat([content, chunk.subarray(0, bytesRead)])
- signal?.throwIfAborted()
- const first = scanZstdFrames(content, 1).frames[0]
- signal?.throwIfAborted()
- if (first === undefined) continue
- let plaintext: Buffer
- try {
- signal?.throwIfAborted()
- plaintext = await decompressZstdFrame(content.subarray(first.start, first.end))
- } catch (error) {
- /* v8 ignore next -- decoder failure plus concurrent abort is timing-dependent */
- if (signal?.aborted) signal.throwIfAborted()
- throw new Error('corrupt Zstandard session log: header frame failed validation', { cause: error })
- }
- signal?.throwIfAborted()
- if (plaintext.length === 0 || plaintext.indexOf(0x0A) !== plaintext.length - 1) {
- throw new Error('corrupt Zstandard session log: first frame is not exactly one header line')
- }
- return plaintext.subarray(0, -1).toString('utf8')
- }
- } finally {
- await handle.close()
- }
- }
- /** Find the unique physical log for an id across every project directory. */
- private async findLog(id: SessionId, signal?: AbortSignal): Promise<string | undefined> {
- const matches: string[] = []
- for (const project of await this.listProjectDirs(signal)) {
- signal?.throwIfAborted()
- await this.rejectLegacyFlatArtifact(project, id, signal)
- signal?.throwIfAborted()
- const dir = join(project, encodeSegment(id))
- const path = join(dir, `session${logSuffix(this.compression)}`)
- const opposite = join(dir, `session${logSuffix(this.oppositeCompression())}`)
- const oppositeExists = await this.exists(opposite)
- signal?.throwIfAborted()
- if (oppositeExists) throw this.encodingMismatch(opposite)
- const pathExists = await this.exists(path)
- signal?.throwIfAborted()
- if (pathExists) matches.push(path)
- }
- if (matches.length > 1) {
- throw new Error(`duplicate JSONL session id "${id}" appears in multiple project directories`)
- }
- signal?.throwIfAborted()
- return matches[0]
- }
- /** Require an existing configured root to be a readable directory. */
- private assertUsableRoot(): void {
- try {
- readdirSync(this.root)
- } catch (error) {
- if (isENOENT(error)) return
- throw error
- }
- }
- /** Reject metadata that does not identify the selected physical log. */
- private async assertStoredIdentity(
- path: string,
- meta: SessionHeader,
- expectedId?: SessionId,
- signal?: AbortSignal,
- ): Promise<void> {
- signal?.throwIfAborted()
- if (expectedId !== undefined && meta.id !== expectedId) {
- throw new Error(`corrupt session log "${path}": requested id "${expectedId}" does not match header id "${meta.id}"`)
- }
- let expectedPath: string
- try {
- expectedPath = logPath(this.root, meta.cwd, meta.id, this.compression)
- } catch (error) {
- throw new Error(`corrupt session log "${path}": header id cannot name a storage path`, { cause: error })
- }
- if (path !== expectedPath && !await this.sameFile(path, expectedPath, signal)) {
- throw new Error(`corrupt session log "${path}": header id "${meta.id}" and cwd identify "${expectedPath}"`)
- }
- signal?.throwIfAborted()
- }
- /**
- * Whether two path spellings resolve to the same physical file. This admits
- * case aliases on case-insensitive filesystems without weakening identity
- * checks on case-sensitive stores.
- */
- private async sameFile(path: string, expectedPath: string, signal?: AbortSignal): Promise<boolean> {
- signal?.throwIfAborted()
- try {
- const [actual, expected] = await Promise.all([realpath(path), realpath(expectedPath)])
- signal?.throwIfAborted()
- return actual === expected
- } catch (error) {
- signal?.throwIfAborted()
- /* v8 ignore else -- non-ENOENT realpath failures require an external permission or I/O fault */
- if (isENOENT(error)) return false
- /* v8 ignore next -- non-ENOENT realpath failures are external I/O faults, propagated unchanged */
- throw error
- }
- }
- /** The human-readable project directories under the configured root. */
- private async listProjectDirs(signal?: AbortSignal): Promise<string[]> {
- try {
- signal?.throwIfAborted()
- const entries = await readdir(this.root, { withFileTypes: true })
- signal?.throwIfAborted()
- return entries.filter(e => e.isDirectory()).map(e => join(this.root, e.name))
- } catch (error) {
- // Only an absent root means no sessions; rethrow every other I/O failure.
- if (isENOENT(error)) return []
- throw error
- }
- }
- /** List session-owned directories and reject the obsolete flat-file layout. */
- private async listSessionDirs(project: string, signal?: AbortSignal): Promise<string[]> {
- signal?.throwIfAborted()
- const entries = await readdir(project, { withFileTypes: true })
- signal?.throwIfAborted()
- const legacy = entries.find(entry =>
- entry.isFile() && (entry.name.endsWith('.jsonl') || entry.name.endsWith('.jsonl.zstd')))
- if (legacy !== undefined) throw this.legacyLayout(join(project, legacy.name))
- return entries.filter(entry => entry.isDirectory()).map(entry => join(project, entry.name))
- }
- /** Reject a root that already belongs to the other physical encoding. */
- private ensureRootEncoding(): Promise<void> {
- this.rootEncodingCheck ??= this.checkRootEncoding()
- return this.rootEncodingCheck
- }
- private async checkRootEncoding(): Promise<void> {
- for (const project of await this.listProjectDirs()) {
- for (const dir of await this.listSessionDirs(project)) {
- const incompatible = join(dir, `session${logSuffix(this.oppositeCompression())}`)
- if (await this.exists(incompatible)) throw this.encodingMismatch(incompatible)
- }
- }
- }
- private async rejectLegacyFlatArtifact(
- project: string,
- id: SessionId,
- signal?: AbortSignal,
- ): Promise<void> {
- signal?.throwIfAborted()
- const encoded = encodeSegment(id)
- for (const compression of ['zstd', 'none'] as const) {
- const path = join(project, encoded + logSuffix(compression))
- const artifactExists = await this.exists(path)
- signal?.throwIfAborted()
- if (artifactExists) throw this.legacyLayout(path)
- }
- }
- private async rejectOppositeArtifact(cwd: string | undefined, id: SessionId): Promise<void> {
- const path = logPath(this.root, cwd, id, this.oppositeCompression())
- if (await this.exists(path)) throw this.encodingMismatch(path)
- }
- private oppositeCompression(): JsonlCompression {
- return this.compression === 'zstd' ? 'none' : 'zstd'
- }
- private encodingMismatch(path: string): Error {
- return new Error(
- `session artifact ${JSON.stringify(path)} uses ${logSuffix(this.oppositeCompression())}, `
- + `but this backend is configured for compression ${JSON.stringify(this.compression)}; `
- + 'use a separate root or select the matching compression mode',
- )
- }
- private legacyLayout(path: string): Error {
- return new Error(
- `session artifact ${JSON.stringify(path)} uses the unsupported flat-file layout; `
- + 'use a separate root or move it into a project/session directory before loading',
- )
- }
- private async exists(path: string): Promise<boolean> {
- try {
- const handle = await open(path, 'r')
- await handle.close()
- return true
- } catch (error) {
- // Only ENOENT means absent. A permission/I/O error must surface rather
- // than letting load or collision checks proceed under false absence.
- // Windows reports ENOENT, not ENOTDIR, for `regular-file/child`; verify
- // the immediate parent so a blocked session directory remains a storage fault.
- /* v8 ignore else -- Windows reports file-valued parents as ENOENT; POSIX covers direct ENOTDIR. */
- if (isENOENT(error)) {
- await this.assertLogParentAllowsAbsence(path)
- return false
- }
- /* v8 ignore next -- Windows repairs ENOTDIR from ENOENT above; POSIX covers direct ENOTDIR. */
- throw error
- }
- }
- /* v8 ignore start -- native Windows coverage exercises this repair; POSIX open reports ENOTDIR before this point. */
- private async assertLogParentAllowsAbsence(path: string): Promise<void> {
- try {
- const parent = dirname(path)
- const info = await stat(parent)
- if (info.isDirectory()) return
- const error = new Error(`ENOTDIR: parent path exists but is not a directory: ${parent}`) as NodeJS.ErrnoException
- error.code = 'ENOTDIR'
- error.path = parent
- throw error
- } catch (error) {
- if (isENOENT(error)) return
- throw error
- }
- }
- /* v8 ignore stop */
- }
- export default SessionPersistenceJsonl
|