返回源码地图

packages/session/session-persistence-jsonl/src/index.ts

main snapshot · da00f7f5358f · 正文引用章节 16 / 16 / 16 / 16;完整原文可核对,不声称全文件人工逐行审计

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

1/**
2 * JSONL durable session-persistence backend. It stores a header and contiguous
3 * events in immutable generation files under one directory per session and serves the handle-based
4 * `SessionPersistence` API: `create`/`open` return per-session handles, and
5 * every read validates the same fail-closed storage contract.
6 * @module @deepseek-ai/dsh-session-persistence-jsonl
7 */
8
9import { Context } from '@deepseek-ai/cordis'
10import z from '@deepseek-ai/schemastery'
11import {
12 createSessionFormatCatalogWithChildren,
13 SessionFormatUnsupportedMigrationError,
14 sessionFormatCatalog,
15} from '@deepseek-ai/dsh-session-format-catalog'
16import { readdirSync, type Dirent } from 'node:fs'
17import { open, mkdir, readdir, realpath, link, rm, stat, truncate } from 'node:fs/promises'
18import { dirname, join, resolve } from 'node:path'
19import { performance } from 'node:perf_hooks'
20import { scheduler } from 'node:timers/promises'
21import { createHash, randomBytes } from 'node:crypto'
22import {
23 SessionPersistence, SessionPersistenceRevision, SessionFormatUnsupportedError,
24 SessionPersistenceCorruptionError,
25 SessionAlreadyExistsError, SessionPersistenceNotFoundError,
26 assertStoredId, materializeCreateHeader, sessionFormatVersionRefusal, validateStoredEvents,
27 type SessionAccess, type SessionHandle,
28 type SessionHandleReadResult,
29 type SessionLocation, type SessionPersistenceCreateOptions,
30 type SessionPersistenceListOptions, type SessionPersistenceOpenOptions,
31 type SessionPersistenceSnapshot, type SessionPersistenceStatOptions,
32 type SessionPersistenceRevision as PersistenceRevision,
33} from '@deepseek-ai/dsh-session-persistence'
34import { JsonlBackendTracker, JsonlSessionHandle, type StorageHandleState } from './storage.ts'
35import { SessionWriteLease } from './lease.ts'
36import { SESSION_FORMAT_VERSION, SessionId as makeSessionId, SessionLogOffset } from '@deepseek-ai/dsh-session'
37import type { SessionEvent, SessionId, SessionHeader, SessionLogOffset as SessionLogOffsetType } from '@deepseek-ai/dsh-session'
38import {
39 assertNoRetiredHeaderFields, encodeSegment, eventLines, generationLogFilename, generationLogPath, logPath, logSuffix,
40 parseGenerationLogFilename, projectDir, scanLog, sessionDir, SessionLogScanner, toHeaderLine,
41 type JsonlCompression,
42} from './format.ts'
43import {
44 compressZstdFrame, createZstdFrameDecoder, decompressZstdFrame, decompressZstdPrefix, scanZstdFrames,
45} from './zstd.ts'
46import { ensureDurableDirectoryWin32, publishNewFileWin32 } from './win32.ts'
47import { verifyCurrentGenerationInWorker } from './migration-verifier.ts'
48import { prepareCatalogFacts } from './catalog-migration.ts'
49import {
50 JsonlGenerationSourceChangedError,
51 JsonlGenerationUnsupportedMigrationError,
52 prepareJsonlMigration,
53 readStableJsonlFile,
54 type JsonlGenerationFormatAdapter,
55 type JsonlPhysicalIdentity,
56 type PreparedJsonlMigration,
57} from './generation.ts'
58
59export type { JsonlCompression } from './format.ts'
60
61/**
62 * Internal handoff-reuse policy, not deployment configuration: a cold
63 * observation and the resume that immediately follows it reuse one parsed
64 * log, so the memo only needs the sessions in flight between those steps.
65 */
66const COLD_LOG_MEMO_MAX_ENTRIES = 2
67
68const DEFAULT_COMPRESSION: JsonlCompression = 'zstd'
69/**
70 * Internal scheduling constant, not deployment configuration: balance
71 * frame-boundary event-loop yields against `setImmediate` overhead. One frame
72 * remains an indivisible synchronous decode.
73 */
74const ZSTD_DECODE_YIELD_INTERVAL_MS = 500
75
76/** Assert that the independently decodable first frame contains only the header record. */
77function assertZstdHeaderFrame(plaintext: Buffer): void {
78 if (plaintext.length === 0 || plaintext.indexOf(0x0A) !== plaintext.length - 1) {
79 throw new Error('corrupt Zstandard session log: first frame is not exactly one header line')
80 }
81}
82
83/** Loader schema for the JSONL artifact's physical encoding. */
84export const JsonlCompressionSchema: z<JsonlCompression> = z.union([
85 z.const('zstd'),
86 z.const('none'),
87]).default(DEFAULT_COMPRESSION)
88
89/** Plugin config for the JSONL backend's root and physical encoding. */
90export interface Config {
91 /**
92 * Root directory for all session files. Required (no default): a default of
93 * `process.cwd()` would scatter session files as the process's cwd changes
94 * (bash calls, subprocesses). Sessions group under human-readable project
95 * directories, then per-session directories. An existing root must be a
96 * readable directory; an absent root is created on first materialization.
97 */
98 root: string
99 /** Physical encoding; defaults to checksummed Zstandard frames. */
100 compression?: JsonlCompression
101}
102
103/** One stored event graph whose producer has established immutable sharing. */
104interface FrozenStoredEvents extends SessionHandleReadResult {
105 readonly eventState: 'shared-frozen'
106}
107
108/** State shared by prepared historical and published current logs. */
109interface StoredLogBase extends FrozenStoredEvents {
110 readonly meta: SessionHeader
111 readonly tornTruncateTo: number | undefined
112 /** Complete events recovered from the torn final frame; the write path rewrites them durably. */
113 readonly recoveredTail: SessionEvent[]
114 /** Exact fork-inherited prefix length stored in the header line. */
115 readonly inheritedEventCount: SessionLogOffsetType
116 readonly revision: PersistenceRevision
117}
118
119/** A decoded current generation that is already durable. */
120interface CurrentStoredLog extends StoredLogBase {
121 readonly status: 'current'
122}
123
124/** A migrated historical generation retained until an explicit write open publishes it. */
125interface PreparedStoredLog extends StoredLogBase {
126 readonly status: 'prepared'
127 readonly validateRelatedSources: () => Promise<void>
128 readonly publication: {
129 readonly source: ResolvedJsonlGeneration
130 readonly value: PreparedJsonlMigration
131 }
132}
133
134/** A validated logical log, either durable current state or prepared historical state. */
135type StoredLog = CurrentStoredLog | PreparedStoredLog
136
137/** Deep-freeze acyclic stored JSON; its arrays contain only indexed JSON values. */
138function freezeStoredEvent(event: SessionEvent): void {
139 const pending: object[] = [event]
140 while (pending.length > 0) {
141 // The non-empty check proves an object remains to visit.
142 // oxlint-disable-next-line typescript/no-non-null-assertion
143 const current = pending.pop()!
144 Object.freeze(current)
145 if (Array.isArray(current)) {
146 for (let index = 0; index < current.length; index += 1) {
147 const child: unknown = current[index]
148 if (child !== null && typeof child === 'object') pending.push(child)
149 }
150 } else {
151 for (const key in current) {
152 const child = (current as Record<string, unknown>)[key]
153 if (child !== null && typeof child === 'object') pending.push(child)
154 }
155 }
156 }
157}
158
159/** Establish immutable sharing for one decoded event graph and report that state. */
160function freezeStoredEvents(events: SessionEvent[]): FrozenStoredEvents {
161 for (const event of events) freezeStoredEvent(event)
162 Object.freeze(events)
163 return { eventState: 'shared-frozen', events }
164}
165
166/** One authoritative immutable generation selected from a Session directory. */
167interface ResolvedJsonlGeneration {
168 readonly sourcePath: string
169 readonly sourceVersion: number
170 readonly currentPath: string
171}
172
173/** One backend-owned historical preparation shared by its current callers. */
174interface MigrationPreparation {
175 readonly sourcePath: string
176 readonly sourceRevision: PersistenceRevision
177 readonly controller: AbortController
178 readonly promise: Promise<PreparedStoredLog>
179 settled: boolean
180 waiters: number
181}
182
183/** Build the stat-derived best-effort change token shared by full and lightweight reads. */
184function fileRevision(identity: JsonlPhysicalIdentity): PersistenceRevision {
185 return SessionPersistenceRevision([
186 identity.dev,
187 identity.ino,
188 identity.size,
189 identity.mtimeNs,
190 identity.ctimeNs,
191 ].join(':'))
192}
193
194/** Whether a filesystem error means absence; every non-ENOENT failure must surface. */
195function isENOENT(error: unknown): boolean {
196 return (error as NodeJS.ErrnoException | null)?.code === 'ENOENT'
197}
198
199/** Whether a filesystem-owned failure should retain its original errno and path. */
200function isErrnoException(error: unknown): error is NodeJS.ErrnoException {
201 return typeof (error as NodeJS.ErrnoException | null)?.code === 'string'
202}
203
204/** Preserve an Error abort reason and normalize hostile non-Error reasons. */
205function abortError(signal: AbortSignal): Error {
206 return signal.reason instanceof Error
207 ? signal.reason
208 : new Error('session migration preparation aborted', { cause: signal.reason })
209}
210
211/** Let one caller stop waiting without transferring cancellation ownership to shared work. */
212function waitWithAbort<T>(operation: Promise<T>, signal?: AbortSignal): Promise<T> {
213 if (signal === undefined) return operation
214 /* v8 ignore next -- requireStoredLog synchronously rechecks the signal immediately before waiting. */
215 if (signal.aborted) return Promise.reject(abortError(signal))
216 return new Promise<T>((resolve, reject) => {
217 const stopWaiting = (): void => {
218 reject(abortError(signal))
219 }
220 signal.addEventListener('abort', stopWaiting, { once: true })
221 void operation.then(
222 (value) => {
223 signal.removeEventListener('abort', stopWaiting)
224 resolve(value)
225 },
226 (error: unknown) => {
227 signal.removeEventListener('abort', stopWaiting)
228 /* v8 ignore else -- the preparation owner normalizes every rejection before this waiter sees it. */
229 if (error instanceof Error) {
230 reject(error)
231 } else {
232 reject(new Error('session migration preparation failed', { cause: error }))
233 }
234 },
235 )
236 })
237}
238
239/**
240 * The JSONL persistence backend. Load as a plugin; it registers as
241 * `ctx.sessionPersistence`. Sessions materialize lazily: a created session is
242 * visible to this process immediately, reaches disk on its first append or
243 * flush, and never existed if the process crashes before that.
244 */
245class JsonlSessionPersistence extends SessionPersistence {
246 static Config: z<Config> = z.object({
247 root: z.string().required(),
248 compression: JsonlCompressionSchema,
249 })
250
251 /** Backend label for diagnostics and effects; shadows `Service.name` without changing the service key. */
252 override readonly name = 'session-persistence-jsonl'
253
254 private root: string
255 private compression: JsonlCompression
256 private rootEncodingCheck: Promise<void> | undefined
257 private readonly tracker = new JsonlBackendTracker(this.name)
258 private readonly generationFormat: Omit<JsonlGenerationFormatAdapter, 'createRestore'>
259 /**
260 * Bounded LRU of parsed, validated stored logs keyed by session id and
261 * guarded by the stat-derived revision, so an immediate cold-read handoff
262 * (observation then resume) parses the artifact once. Every local mutation
263 * for an id invalidates its entry; a foreign write misses through the
264 * revision guard.
265 */
266 private readonly coldLogMemo = new Map<SessionId, StoredLog>()
267 /** One joinable decode/migration operation per selected historical Session file revision. */
268 private readonly migrationPreparations = new Map<SessionId, MigrationPreparation>()
269
270 constructor(ctx: Context, public config: Config) {
271 super(ctx)
272 /* v8 ignore next 5 -- generated catalog and Session source share one build-time version owner. */
273 if (sessionFormatCatalog.currentVersion !== SESSION_FORMAT_VERSION) {
274 throw new Error(
275 `session-persistence-jsonl: format catalog v${sessionFormatCatalog.currentVersion} `
276 + `does not match Session v${SESSION_FORMAT_VERSION}`,
277 )
278 }
279 // Resolve once so later process.cwd() changes cannot split one backend across roots.
280 this.root = resolve(config.root)
281 this.compression = config.compression ?? DEFAULT_COMPRESSION
282 this.generationFormat = {
283 currentVersion: sessionFormatCatalog.currentVersion,
284 encodeHeader: (header, inheritedEventCount) =>
285 sessionFormatCatalog.encodeCurrentHeader(header, inheritedEventCount),
286 encodeEvent: event => sessionFormatCatalog.encodeCurrentEvent(event),
287 isUnsupportedMigrationError: (error): error is SessionFormatUnsupportedMigrationError =>
288 error instanceof SessionFormatUnsupportedMigrationError,
289 }
290 this.assertUsableRoot()
291 this.tracker.install(ctx)
292 }
293
294 /**
295 * Refusal-diagnostics hook: the absolute target path, without touching the filesystem.
296 * @param meta - the stored header naming the session and its cwd.
297 * @returns the artifact kind and absolute path.
298 */
299 private locate(meta: SessionHeader): SessionLocation {
300 return { kind: 'jsonl', path: logPath(this.root, meta.cwd, meta.id, this.compression) }
301 }
302
303 // --- SessionPersistence service API ---
304
305 /**
306 * Create a new stored session and take its write ownership. The session is
307 * visible to this process immediately; the physical artifact appears on the
308 * first append or flush.
309 * @param header - the immutable header to store; must be losslessly
310 * JSON-serializable with a non-negative safe-integer `createdAt`.
311 * @param options - optional cancellation.
312 * @returns the owned write handle.
313 */
314 async create(header: SessionHeader, options?: SessionPersistenceCreateOptions): Promise<SessionHandle> {
315 options?.signal?.throwIfAborted()
316 const snapshot = materializeCreateHeader(header)
317 // Fail fast on a seeded/cut mismatch with the exact refusal the header
318 // line encoder enforces at materialization.
319 toHeaderLine(snapshot, options?.inheritedEventCount)
320 const inheritedEventCount = SessionLogOffset(options?.inheritedEventCount ?? 0)
321 await this.ensureRootEncoding()
322 options?.signal?.throwIfAborted()
323 if (this.tracker.hasPending(snapshot.id) || await this.findLog(snapshot.id, options?.signal) !== undefined) {
324 throw new SessionAlreadyExistsError(snapshot.id)
325 }
326 options?.signal?.throwIfAborted()
327 // No lock yet: before materialization there is no durable artifact for
328 // another process to contend over, so the handle acquires the lock right
329 // before its first log bytes publish (ensureLease); an unmaterialized
330 // session leaves no filesystem footprint at all.
331 this.tracker.registerCreated(snapshot, inheritedEventCount)
332 return this.tracker.adopt(new JsonlSessionHandle(this, snapshot.id, snapshot, 'write', { cursor: 0, materialized: false, inheritedEventCount }))
333 }
334
335 /**
336 * Open an existing stored session for `read` or single-writer `write`.
337 * @param id - the stored session to open.
338 * @param access - `read` (no ownership) or `write` (atomic in-process claim).
339 * @param options - optional cancellation.
340 * @returns the open handle.
341 */
342 async open(id: SessionId, access: SessionAccess, options?: SessionPersistenceOpenOptions): Promise<SessionHandle> {
343 options?.signal?.throwIfAborted()
344 await this.ensureRootEncoding()
345 options?.signal?.throwIfAborted()
346 const pending = this.tracker.pendingOf(id)
347 if (access === 'read') {
348 if (pending !== undefined) {
349 return this.tracker.adopt(new JsonlSessionHandle(this, id, pending.header, 'read', { cursor: 0, materialized: false, inheritedEventCount: pending.inheritedEventCount }))
350 }
351 let stored: StoredLog
352 try {
353 stored = await this.requireStoredLog(id, options?.signal)
354 } catch (error: unknown) {
355 if (!(error instanceof JsonlGenerationSourceChangedError)) throw error
356 stored = await this.requireStoredLog(id, options?.signal)
357 }
358 let state: StorageHandleState
359 if (stored.status === 'prepared') {
360 state = {
361 cursor: 0,
362 materialized: true,
363 inheritedEventCount: stored.inheritedEventCount,
364 primed: stored,
365 }
366 } else {
367 state = {
368 cursor: 0,
369 materialized: true,
370 inheritedEventCount: stored.inheritedEventCount,
371 }
372 }
373 return this.tracker.adopt(new JsonlSessionHandle(this, id, stored.meta, 'read', state))
374 }
375 // A pending entry always belongs to an ACTIVE creator handle (close erases
376 // it), so the claim below rejects that case as already owned.
377 this.tracker.claimWrite(id)
378 let lease: SessionWriteLease | undefined
379 try {
380 const resolved = await this.findLog(id, options?.signal)
381 if (resolved === undefined) throw new SessionPersistenceNotFoundError(id)
382 lease = await this.acquireLease(id, undefined, dirname(resolved.currentPath))
383 const prepared = await this.requireStoredLog(id, options?.signal)
384 options?.signal?.throwIfAborted()
385 let stored: CurrentStoredLog
386 if (prepared.status === 'prepared') {
387 stored = await this.publishStoredMigration(id, prepared)
388 } else {
389 stored = prepared
390 }
391 options?.signal?.throwIfAborted()
392 return this.tracker.adopt(new JsonlSessionHandle(this, id, stored.meta, 'write', {
393 cursor: stored.events.length,
394 materialized: true,
395 tornTruncateTo: stored.tornTruncateTo,
396 recoveredTail: stored.recoveredTail,
397 inheritedEventCount: stored.inheritedEventCount,
398 primed: stored,
399 }, lease))
400 } catch (error) {
401 // Free the in-process claim no matter how the kernel-lock release
402 // fares, and keep the original diagnostic: a release failure joins it
403 // instead of replacing it.
404 /* v8 ignore next -- typed backends and fs reject with Error */
405 const failure = error instanceof Error ? error : new Error(String(error))
406 let releaseFailure: Error | undefined
407 try {
408 await lease?.release()
409 } catch (raw: unknown) {
410 /* v8 ignore next -- lock releases reject with Error */
411 releaseFailure = raw instanceof Error ? raw : new Error(String(raw))
412 }
413 this.tracker.releaseClaim(id)
414 if (releaseFailure !== undefined) {
415 throw new AggregateError([failure, releaseFailure], `session "${id}": write open failed and its lock release failed`)
416 }
417 throw failure
418 }
419 }
420
421 /**
422 * Flush every active write handle in one durability barrier; see the seam
423 * contract.
424 * @returns resolution once every write handle active at the call has flushed.
425 */
426 flush(): Promise<void> {
427 return this.tracker.flushAll()
428 }
429
430 /**
431 * Observe one stored session without reading its event log.
432 * @param id - the stored session to observe.
433 * @param options - optional cancellation.
434 * @returns the snapshot (`sizeBytes` carries the physical artifact size), or
435 * `undefined` when the session does not exist.
436 */
437 async stat(
438 id: SessionId,
439 options?: SessionPersistenceStatOptions,
440 ): Promise<SessionPersistenceSnapshot | undefined> {
441 options?.signal?.throwIfAborted()
442 await this.ensureRootEncoding()
443 options?.signal?.throwIfAborted()
444 const pending = this.tracker.pendingOf(id)
445 if (pending !== undefined) {
446 return { header: pending.header, revision: pending.revision }
447 }
448 const selected = await this.findLog(id, options?.signal)
449 if (selected === undefined) return undefined
450 const header = await this.readGenerationHeader(selected, id, options?.signal)
451 if (header === undefined) return undefined
452 try {
453 const identity = await stat(selected.sourcePath, { bigint: true })
454 options?.signal?.throwIfAborted()
455 return {
456 header,
457 revision: selected.sourceVersion < SESSION_FORMAT_VERSION
458 ? SessionPersistenceRevision(`${fileRevision(identity)}:${await this.historicalCorpusRevision(options?.signal)}`)
459 : fileRevision(identity),
460 sizeBytes: Number(identity.size),
461 }
462 } catch (error: unknown) {
463 options?.signal?.throwIfAborted()
464 if (isENOENT(error)) return undefined
465 throw error
466 }
467 }
468
469 /**
470 * List every stored session visible to this process: materialized artifacts
471 * plus this process's created-but-unmaterialized sessions.
472 * @param options - optional cancellation.
473 * @returns one snapshot per session, in no promised order.
474 */
475 async list(options?: SessionPersistenceListOptions): Promise<readonly SessionPersistenceSnapshot[]> {
476 const signal = options?.signal
477 const snapshots: SessionPersistenceSnapshot[] = []
478 const listed = new Set<SessionId>()
479 // Snapshot pending entries BEFORE scanning storage: a session whose first
480 // append lands mid-scan is then still in this snapshot (its artifact may
481 // predate the scan), so create-to-list visibility never has a hole.
482 const pending = [...this.tracker.pendingEntries()]
483 const artifacts = await this.listArtifacts(signal)
484 const corpusRevision = artifacts.some(artifact => artifact.sourceVersion < SESSION_FORMAT_VERSION)
485 ? await this.historicalCorpusRevision(signal) : undefined
486 for (const artifact of artifacts) {
487 signal?.throwIfAborted()
488 try {
489 const identity = await stat(artifact.path, { bigint: true })
490 signal?.throwIfAborted()
491 listed.add(artifact.header.id)
492 snapshots.push({
493 header: artifact.header,
494 revision: artifact.sourceVersion < SESSION_FORMAT_VERSION
495 ? SessionPersistenceRevision(`${fileRevision(identity)}:${corpusRevision}`)
496 : fileRevision(identity),
497 sizeBytes: Number(identity.size),
498 })
499 } catch (error: unknown) {
500 signal?.throwIfAborted()
501 if (!isENOENT(error)) throw error
502 }
503 }
504 for (const [id, entry] of pending) {
505 if (!listed.has(id)) snapshots.push({ header: entry.header, revision: entry.revision })
506 }
507 signal?.throwIfAborted()
508 return snapshots
509 }
510
511 // --- handle-facing storage internals (package-private via the handle class below) ---
512
513 /** Resolve and read one stored log, refusing loudly when the artifact is absent. */
514 private async requireStoredLog(id: SessionId, signal?: AbortSignal): Promise<StoredLog> {
515 const selected = await this.findLog(id, signal)
516 if (selected === undefined) throw new SessionPersistenceNotFoundError(id)
517 if (selected.sourceVersion < SESSION_FORMAT_VERSION) {
518 const sourceRevision = fileRevision(await stat(selected.sourcePath, { bigint: true }))
519 signal?.throwIfAborted()
520 let preparation = this.migrationPreparations.get(id)
521 if (preparation === undefined
522 || preparation.sourcePath !== selected.sourcePath
523 || preparation.sourceRevision !== sourceRevision) {
524 const controller = new AbortController()
525 const promise = this.loadStoredMigration(id, selected, sourceRevision, controller.signal)
526 preparation = {
527 sourcePath: selected.sourcePath,
528 sourceRevision,
529 controller,
530 promise,
531 settled: false,
532 waiters: 0,
533 }
534 this.migrationPreparations.set(id, preparation)
535 const created = preparation
536 const release = (): void => {
537 created.settled = true
538 if (this.migrationPreparations.get(id) === created) {
539 this.migrationPreparations.delete(id)
540 }
541 }
542 void promise.then(release, release)
543 }
544 signal?.throwIfAborted()
545 return this.waitForPreparation(id, preparation, signal)
546 }
547 if (selected.sourceVersion > SESSION_FORMAT_VERSION) {
548 const header = await this.readGenerationHeader(selected, id, signal)
549 /* v8 ignore else -- a readable future header is rejected inside readGenerationHeader. */
550 if (header === undefined) {
551 throw new SessionPersistenceCorruptionError(
552 `session "${id}": stored log has a malformed header (raw log: ${selected.sourcePath})`,
553 { cause: new Error('malformed Session header') },
554 )
555 }
556 /* v8 ignore next -- readGenerationHeader rejects every future version. */
557 throw new SessionFormatUnsupportedError(
558 `${sessionFormatVersionRefusal(id, selected.sourceVersion)} (raw log: ${selected.sourcePath})`,
559 { kind: 'jsonl', path: selected.sourcePath },
560 )
561 }
562 const probe = fileRevision(await stat(selected.sourcePath, { bigint: true }))
563 const memoized = this.coldLogMemo.get(id)
564 if (memoized?.status === 'current' && memoized.revision === probe) {
565 this.coldLogMemo.delete(id)
566 this.coldLogMemo.set(id, memoized)
567 return memoized
568 }
569 const current = await readStableJsonlFile(selected.sourcePath, signal)
570 return this.decodeStoredLog(
571 selected.sourcePath,
572 id,
573 current.bytes,
574 fileRevision(current.identity),
575 signal,
576 )
577 }
578
579 /** Probe the memo and otherwise decode one historical generation under backend cancellation. */
580 private async loadStoredMigration(
581 id: SessionId,
582 selected: ResolvedJsonlGeneration,
583 sourceRevision: PersistenceRevision,
584 signal: AbortSignal,
585 ): Promise<PreparedStoredLog> {
586 signal.throwIfAborted()
587 const memoized = this.coldLogMemo.get(id)
588 if (memoized?.status === 'prepared' && memoized.revision === sourceRevision) {
589 try {
590 await memoized.validateRelatedSources()
591 } catch (error: unknown) {
592 this.coldLogMemo.delete(id)
593 throw this.generationFailure(id, selected, error)
594 }
595 this.coldLogMemo.delete(id)
596 this.coldLogMemo.set(id, memoized)
597 return memoized
598 }
599 return this.prepareStoredMigration(id, selected, signal)
600 }
601
602 /** Await shared preparation for one caller and abort it only after its last waiter leaves. */
603 private async waitForPreparation(
604 id: SessionId,
605 preparation: MigrationPreparation,
606 signal?: AbortSignal,
607 ): Promise<PreparedStoredLog> {
608 preparation.waiters += 1
609 try {
610 return await waitWithAbort(preparation.promise, signal)
611 } finally {
612 preparation.waiters -= 1
613 if (preparation.waiters === 0 && !preparation.settled) {
614 /* v8 ignore else -- a newer selected source may already own this id's preparation slot. */
615 if (this.migrationPreparations.get(id) === preparation) {
616 this.migrationPreparations.delete(id)
617 }
618 preparation.controller.abort()
619 }
620 }
621 }
622
623 /** Decode one historical generation without publishing a successor. */
624 private async prepareStoredMigration(
625 id: SessionId,
626 selected: ResolvedJsonlGeneration,
627 signal: AbortSignal,
628 ): Promise<PreparedStoredLog> {
629 let prepared: Awaited<ReturnType<typeof prepareJsonlMigration>>
630 let validateRelatedSources: () => Promise<void>
631 try {
632 const children = async () => (await this.listArtifacts(signal))
633 .filter(source => source.header.origin === 'subagent' && source.header.parentSession === id)
634 const sources = await children()
635 const related = await prepareCatalogFacts(id, sources, this.compression, signal)
636 for (const failure of related.failures) {
637 this.ctx.logger.warn(`${this.name}: session "${id}" catalog retained a child with unknown descriptor (raw log: ${failure.path}): ${String(failure.error)}`)
638 }
639 const membership = sources.map(source => source.path).sort()
640 validateRelatedSources = async () => {
641 const current = (await children()).map(source => source.path).sort()
642 const before = new Set(membership)
643 const after = new Set(current)
644 const changed = current.find(path => !before.has(path)) ?? membership.find(path => !after.has(path))
645 if (changed !== undefined) throw new JsonlGenerationSourceChangedError(changed)
646 await related.validate()
647 }
648 prepared = await prepareJsonlMigration({
649 sourcePath: selected.sourcePath,
650 sourceVersion: selected.sourceVersion,
651 currentPath: selected.currentPath,
652 compression: this.compression,
653 format: {
654 ...this.generationFormat,
655 createRestore: header => createSessionFormatCatalogWithChildren(related.facts).createRestore(header, {
656 recovery: 'recoverable', validation: 'transformed',
657 }),
658 },
659 validateRelatedSources,
660 verifyCurrentFile: verifyCurrentGenerationInWorker,
661 validateHistoricalHeader: headerValue => this.validateSourceIdentity(
662 selected,
663 headerValue,
664 id,
665 signal,
666 ),
667 signal,
668 })
669 } catch (error: unknown) {
670 throw this.generationFailure(id, selected, error)
671 }
672 const meta = this.currentHeader(prepared.artifact.header)
673 assertStoredId(id, meta)
674 const events = prepared.artifact.events as SessionEvent[]
675 validateStoredEvents(meta, events, { kind: 'jsonl', path: selected.sourcePath })
676 const stored: PreparedStoredLog = {
677 status: 'prepared',
678 validateRelatedSources,
679 meta,
680 ...freezeStoredEvents(events),
681 tornTruncateTo: undefined,
682 recoveredTail: [],
683 inheritedEventCount: SessionLogOffset(prepared.artifact.inheritedEventCount),
684 revision: fileRevision(prepared.sourceIdentity),
685 publication: { source: selected, value: prepared },
686 }
687 this.memoizeStoredLog(id, stored)
688 return stored
689 }
690
691 /** Publish a prepared historical log before granting write access. */
692 private async publishStoredMigration(id: SessionId, stored: PreparedStoredLog): Promise<CurrentStoredLog> {
693 const migration = stored.publication
694 let identity: JsonlPhysicalIdentity
695 try {
696 identity = await migration.value.publish()
697 } catch (error: unknown) {
698 /* v8 ignore else -- a newer preparation may have replaced this stale cache entry. */
699 if (this.coldLogMemo.get(id) === stored) this.coldLogMemo.delete(id)
700 throw this.generationFailure(id, migration.source, error)
701 }
702 const published: CurrentStoredLog = {
703 status: 'current',
704 meta: stored.meta,
705 eventState: stored.eventState,
706 events: stored.events,
707 tornTruncateTo: stored.tornTruncateTo,
708 recoveredTail: stored.recoveredTail,
709 inheritedEventCount: stored.inheritedEventCount,
710 revision: fileRevision(identity),
711 }
712 this.memoizeStoredLog(id, published)
713 return published
714 }
715
716 /** Translate generation-layer failures into the persistence seam's error vocabulary. */
717 private generationFailure(
718 id: SessionId,
719 selected: ResolvedJsonlGeneration,
720 error: unknown,
721 ): Error {
722 if (error instanceof JsonlGenerationUnsupportedMigrationError) {
723 return new SessionFormatUnsupportedError(
724 `${error.message}; source v${error.fromVersion} artifact remains unchanged (raw log: ${selected.sourcePath})`,
725 { kind: 'jsonl', path: selected.sourcePath },
726 )
727 }
728 if (error instanceof JsonlGenerationSourceChangedError) return error
729 if (error instanceof SessionFormatUnsupportedError
730 || error instanceof SessionPersistenceCorruptionError
731 || isErrnoException(error)
732 || error instanceof DOMException && error.name === 'AbortError') return error
733 return new SessionPersistenceCorruptionError(
734 `session "${id}": stored log is corrupt: ${String(error)} (raw log: ${selected.sourcePath})`,
735 { cause: error },
736 )
737 }
738
739 /**
740 * Read, parse, and validate one stored log as the current logical prefix.
741 * @param path - the artifact file to read.
742 * @param expectedId - the session identity the artifact must carry.
743 * @param signal - optional cancellation for the stat/read/decode work.
744 * @returns the validated stored log with any torn-tail truncation point.
745 */
746 async readStoredLog(path: string, expectedId: SessionId, signal?: AbortSignal): Promise<CurrentStoredLog> {
747 signal?.throwIfAborted()
748 const probe = fileRevision(await stat(path, { bigint: true }))
749 const memoized = this.coldLogMemo.get(expectedId)
750 if (memoized?.status === 'current' && memoized.revision === probe) {
751 this.coldLogMemo.delete(expectedId)
752 this.coldLogMemo.set(expectedId, memoized)
753 return memoized
754 }
755 const { bytes, identity } = await readStableJsonlFile(path, signal)
756 return this.decodeStoredLog(path, expectedId, bytes, fileRevision(identity), signal)
757 }
758
759 /** Decode and memoize one already-stable current physical snapshot. */
760 private async decodeStoredLog(
761 path: string,
762 expectedId: SessionId,
763 buffer: Buffer,
764 revision: PersistenceRevision,
765 signal?: AbortSignal,
766 ): Promise<CurrentStoredLog> {
767 let parsed: {
768 meta: SessionHeader
769 inheritedEventCount: SessionLogOffsetType
770 events: SessionEvent[]
771 tornTruncateTo: number | undefined
772 recoveredTail: SessionEvent[]
773 }
774 try {
775 if (this.compression === 'zstd') {
776 parsed = await this.readZstdPrefix(buffer, signal)
777 } else {
778 signal?.throwIfAborted()
779 const { meta, inheritedEventCount, events, committedBytes } = scanLog(buffer)
780 signal?.throwIfAborted()
781 parsed = {
782 meta,
783 inheritedEventCount,
784 events,
785 tornTruncateTo: committedBytes < buffer.byteLength ? committedBytes : undefined,
786 // A torn raw tail is one incomplete JSONL line; it holds no complete
787 // record to recover.
788 recoveredTail: [],
789 }
790 }
791 } catch (error: unknown) {
792 signal?.throwIfAborted()
793 // A parse-time format refusal predates any SessionHeader, so attach the
794 // artifact this read actually refused; every other parse failure is
795 // committed bytes the decoder cannot interpret — damage, classified for
796 // the seam's stable error vocabulary.
797 if (error instanceof SessionFormatUnsupportedError) {
798 throw new SessionFormatUnsupportedError(`${error.message} (raw log: ${path})`, { kind: 'jsonl', path })
799 }
800 throw new SessionPersistenceCorruptionError(`session "${expectedId}": stored log is corrupt: ${String(error)} (raw log: ${path})`, { cause: error })
801 }
802 signal?.throwIfAborted()
803 await this.assertStoredIdentity(path, SESSION_FORMAT_VERSION, parsed.meta, expectedId, signal)
804 signal?.throwIfAborted()
805 assertStoredId(expectedId, parsed.meta)
806 const location = this.locate(parsed.meta)
807 validateStoredEvents(parsed.meta, parsed.events, location)
808 const { events, ...rest } = parsed
809 const stored: CurrentStoredLog = {
810 status: 'current',
811 ...rest,
812 ...freezeStoredEvents(events),
813 revision,
814 }
815 this.memoizeStoredLog(expectedId, stored)
816 return stored
817 }
818
819 /** Insert one parsed log into the bounded handoff cache. */
820 private memoizeStoredLog(id: SessionId, stored: StoredLog): void {
821 this.coldLogMemo.delete(id)
822 this.coldLogMemo.set(id, stored)
823 for (const oldest of this.coldLogMemo.keys()) {
824 if (this.coldLogMemo.size <= COLD_LOG_MEMO_MAX_ENTRIES) break
825 this.coldLogMemo.delete(oldest)
826 }
827 }
828
829 /**
830 * Resolve a session's current-generation log path.
831 * @param id - the stored session to locate.
832 * @param signal - optional cancellation for the directory scans.
833 * @returns the current artifact path, or `undefined` while only a historical generation exists.
834 */
835 async resolveCurrentLog(id: SessionId, signal?: AbortSignal): Promise<string | undefined> {
836 await this.ensureRootEncoding()
837 signal?.throwIfAborted()
838 const selected = await this.findLog(id, signal)
839 if (selected === undefined) return undefined
840 if (selected.sourceVersion === SESSION_FORMAT_VERSION) return selected.sourcePath
841 if (selected.sourceVersion < SESSION_FORMAT_VERSION) return undefined
842 const reason = sessionFormatVersionRefusal(id, selected.sourceVersion)
843 throw new SessionFormatUnsupportedError(
844 `${reason} (raw log: ${selected.sourcePath})`,
845 { kind: 'jsonl', path: selected.sourcePath },
846 )
847 }
848
849 /**
850 * Durably append one validated batch; lazily materializes on the first write.
851 * @param header - the session's stored header.
852 * @param events - the validated contiguous batch, in seq order.
853 * @param isMaterialized - whether the session already has a durable artifact.
854 * @param inheritedEventCount - the exact fork-inherited prefix length written into a materializing header line.
855 */
856 async persistBatch(
857 header: SessionHeader,
858 events: readonly SessionEvent[],
859 isMaterialized: boolean,
860 inheritedEventCount: SessionLogOffsetType,
861 ): Promise<void> {
862 this.coldLogMemo.delete(header.id)
863 await this.ensureRootEncoding()
864 if (isMaterialized) {
865 await this.appendLines(header, events)
866 } else {
867 await this.materialize(header, inheritedEventCount, events)
868 this.tracker.materialized(header.id)
869 }
870 }
871
872 /**
873 * Materialize a header-only artifact for an explicitly durable empty session.
874 * @param header - the session's stored header.
875 * @param inheritedEventCount - the exact fork-inherited prefix length written into the header line.
876 */
877 async persistHeader(header: SessionHeader, inheritedEventCount: SessionLogOffsetType): Promise<void> {
878 this.coldLogMemo.delete(header.id)
879 await this.ensureRootEncoding()
880 await this.materialize(header, inheritedEventCount, [])
881 this.tracker.materialized(header.id)
882 }
883
884 /**
885 * Truncate a torn physical tail durably before this session's first new append.
886 * @param header - the session's stored header.
887 * @param truncateTo - the byte offset the artifact is truncated to.
888 */
889 async truncateTornTail(header: SessionHeader, truncateTo: number): Promise<void> {
890 this.coldLogMemo.delete(header.id)
891 await this.repair(header, truncateTo)
892 this.ctx.logger.warn(`${this.name}: session "${header.id}" recovered from a torn tail; incomplete tail bytes were discarded`)
893 }
894
895 /**
896 * Whether this process still tracks a created-but-unmaterialized session.
897 * @param id - the session to test.
898 * @returns true while the pending entry exists.
899 */
900 hasPendingSession(id: SessionId): boolean {
901 return this.tracker.hasPending(id)
902 }
903
904 /**
905 * Release one handle's backend bookkeeping on close.
906 * @param handle - the closing handle.
907 * @param materialized - whether the session reached durable storage.
908 */
909 releaseHandle(handle: JsonlSessionHandle, materialized: boolean): void {
910 this.tracker.release(handle, materialized)
911 }
912
913 /**
914 * Acquire the session directory's kernel write lock; the kernel holds it
915 * until the handle's close releases the descriptor, including on process death.
916 * @param id - the session the lock guards.
917 * @param cwd - header cwd used to derive the directory for a fresh session.
918 * @param dir - the resolved directory of an existing artifact, when known.
919 * @returns the held lock.
920 */
921 private acquireLease(id: SessionId, cwd: string | undefined, dir = sessionDir(this.root, cwd, id)): Promise<SessionWriteLease> {
922 return SessionWriteLease.acquire(dir, id)
923 }
924
925 /**
926 * Acquire the cross-process write lock for a materializing created session,
927 * called by its handle immediately before the first log bytes publish.
928 * @param header - the session's stored header (its cwd derives the directory).
929 * @returns the held lock.
930 */
931 async acquireWriteLease(header: SessionHeader): Promise<SessionWriteLease> {
932 // Refuse an opposite-encoding artifact before the lock's mkdir publishes
933 // the session directory — the last moment the directory can be absent.
934 await this.rejectOppositeArtifact(header.cwd, header.id)
935 return this.acquireLease(header.id, header.cwd)
936 }
937
938 /** Decode complete frames and retain complete JSONL records from a torn final frame. */
939 private async readZstdPrefix(
940 buffer: Buffer,
941 signal?: AbortSignal,
942 ): Promise<{
943 meta: SessionHeader
944 inheritedEventCount: SessionLogOffsetType
945 events: SessionEvent[]
946 tornTruncateTo: number | undefined
947 recoveredTail: SessionEvent[]
948 }> {
949 signal?.throwIfAborted()
950 const { frames, tornStart } = scanZstdFrames(buffer)
951 signal?.throwIfAborted()
952 if (frames.length === 0) throw new Error('empty or header-less Zstandard session log')
953
954 const decoder = createZstdFrameDecoder()
955 let yieldDeadline = performance.now() + ZSTD_DECODE_YIELD_INTERVAL_MS
956 try {
957 const decodedFrames = decoder.decode(buffer, frames)
958 signal?.throwIfAborted()
959 const headerFrame = decodedFrames.next()
960 signal?.throwIfAborted()
961 /* v8 ignore next -- a non-empty structural frame list makes the decoder yield its first frame or throw. */
962 if (headerFrame.done) throw new Error('empty or header-less Zstandard session log')
963 assertZstdHeaderFrame(headerFrame.value)
964 const scanner = new SessionLogScanner(headerFrame.value)
965
966 let remainingFrames = frames.length - 1
967 for (const plaintext of decodedFrames) {
968 signal?.throwIfAborted()
969 scanner.write(plaintext)
970 remainingFrames -= 1
971 if (remainingFrames > 0 && performance.now() >= yieldDeadline) {
972 await scheduler.yield()
973 signal?.throwIfAborted()
974 yieldDeadline = performance.now() + ZSTD_DECODE_YIELD_INTERVAL_MS
975 }
976 }
977 signal?.throwIfAborted()
978 const complete = scanner.checkpoint()
979 if (complete.committedBytes !== complete.inputBytes) {
980 throw new Error('corrupt Zstandard session log: complete frame contains a torn JSONL record')
981 }
982 if (tornStart === undefined) {
983 const prefix = scanner.finish()
984 return {
985 meta: prefix.meta,
986 inheritedEventCount: prefix.inheritedEventCount,
987 events: prefix.events,
988 tornTruncateTo: undefined,
989 recoveredTail: [],
990 }
991 }
992 // A torn final frame's append never resolved, but complete JSONL records
993 // already flushed into it are real emitted events: recover them, and let
994 // the write path truncate the torn bytes and rewrite them durably.
995 let recoveredPlaintext: Buffer = Buffer.alloc(0)
996 try {
997 signal?.throwIfAborted()
998 recoveredPlaintext = await decompressZstdPrefix(buffer.subarray(tornStart))
999 } catch {
1000 /* v8 ignore next -- decoder failure plus concurrent abort is timing-dependent */
1001 if (signal?.aborted) signal.throwIfAborted()
1002 // A structurally incomplete final frame may end before Node's decoder
1003 // can emit any plaintext; the complete prior frames remain recoverable.
1004 }
1005 signal?.throwIfAborted()
1006 scanner.write(recoveredPlaintext)
1007 const prefix = scanner.finish()
1008 return {
1009 meta: prefix.meta,
1010 inheritedEventCount: prefix.inheritedEventCount,
1011 events: prefix.events,
1012 tornTruncateTo: tornStart,
1013 recoveredTail: prefix.events.slice(complete.eventCount),
1014 }
1015 } catch (error) {
1016 /* v8 ignore next -- decoder failure plus concurrent abort is timing-dependent */
1017 if (signal?.aborted) signal.throwIfAborted()
1018 throw error
1019 } finally {
1020 decoder.close()
1021 }
1022 }
1023
1024 /** Enumerate selected physical generations without interpreting their headers or bodies. */
1025 private async listGenerations(signal?: AbortSignal): Promise<ResolvedJsonlGeneration[]> {
1026 const sources: ResolvedJsonlGeneration[] = []
1027 for (const project of await this.listProjectDirs(signal)) {
1028 for (const dir of await this.listSessionDirs(project, signal)) {
1029 signal?.throwIfAborted()
1030 const selected = await this.resolveGenerationInDirectory(dir, signal)
1031 if (selected !== undefined) sources.push(selected)
1032 }
1033 }
1034 return sources
1035 }
1036
1037 /** Historical logical events depend on the corpus, including members with unreadable headers. */
1038 private async historicalCorpusRevision(signal?: AbortSignal): Promise<string> {
1039 const paths = (await this.listGenerations(signal)).map(source => source.sourcePath).sort()
1040 const hash = createHash('sha256')
1041 for (const path of paths) {
1042 signal?.throwIfAborted()
1043 let revision: string
1044 try {
1045 revision = fileRevision(await stat(path, { bigint: true }))
1046 } catch (error: unknown) {
1047 if (!isENOENT(error)) throw error
1048 revision = 'missing'
1049 }
1050 hash.update(JSON.stringify([path, revision]))
1051 }
1052 signal?.throwIfAborted()
1053 return hash.digest('hex')
1054 }
1055
1056 private async listArtifacts(
1057 signal?: AbortSignal,
1058 ): Promise<Array<{ header: SessionHeader; path: string; sourceVersion: number }>> {
1059 signal?.throwIfAborted()
1060 await this.ensureRootEncoding()
1061 signal?.throwIfAborted()
1062 const artifacts: Array<{ header: SessionHeader; path: string; sourceVersion: number }> = []
1063 const ids = new Set<SessionId>()
1064 for (const selected of await this.listGenerations(signal)) {
1065 signal?.throwIfAborted()
1066 let header: SessionHeader | undefined
1067 try {
1068 header = await this.readGenerationHeader(selected, undefined, signal)
1069 } catch (error: unknown) {
1070 if (error instanceof SessionFormatUnsupportedError || error instanceof SessionPersistenceCorruptionError) continue
1071 throw error
1072 }
1073 if (header === undefined) {
1074 continue
1075 }
1076 if (ids.has(header.id)) {
1077 throw new Error(`duplicate JSONL session id "${header.id}" appears in multiple project directories`)
1078 }
1079 ids.add(header.id)
1080 artifacts.push({ header, path: selected.sourcePath, sourceVersion: selected.sourceVersion })
1081 }
1082 signal?.throwIfAborted()
1083 return artifacts
1084 }
1085
1086 /** Read and translate one selected generation header without inspecting its body. */
1087 private async readGenerationHeader(
1088 selected: ResolvedJsonlGeneration,
1089 expectedId?: SessionId,
1090 signal?: AbortSignal,
1091 ): Promise<SessionHeader | undefined> {
1092 let first: string | undefined
1093 try {
1094 first = this.compression === 'zstd'
1095 ? await this.readFirstZstdLine(selected.sourcePath, signal)
1096 : await this.readFirstLine(selected.sourcePath, signal)
1097 } catch (error: unknown) {
1098 signal?.throwIfAborted()
1099 if (isENOENT(error)) return undefined
1100 throw error
1101 }
1102 signal?.throwIfAborted()
1103 if (first === undefined) return undefined
1104 let value: unknown
1105 try {
1106 value = JSON.parse(first)
1107 } catch {
1108 return undefined
1109 }
1110 assertNoRetiredHeaderFields(value)
1111 const result = sessionFormatCatalog.readHeader(value)
1112 if ('storedVersion' in result && result.storedVersion !== selected.sourceVersion) {
1113 throw new Error(
1114 `session generation filename identifies v${selected.sourceVersion}, `
1115 + `but its header identifies v${result.storedVersion}`,
1116 )
1117 }
1118 if (result.status === 'unsupported') {
1119 const physicalId = String((value as { id?: unknown }).id)
1120 let reason = result.reason
1121 /* v8 ignore else -- released historical header migrations cannot refuse after physical decoding. */
1122 if (result.storedVersion > SESSION_FORMAT_VERSION) {
1123 reason = sessionFormatVersionRefusal(physicalId, result.storedVersion)
1124 }
1125 throw new SessionFormatUnsupportedError(
1126 `${reason} (raw log: ${selected.sourcePath})`,
1127 { kind: 'jsonl', path: selected.sourcePath },
1128 )
1129 }
1130 if (result.status === 'malformed') return undefined
1131 const header = this.currentHeader(result.header)
1132 await this.assertStoredIdentity(
1133 selected.sourcePath,
1134 selected.sourceVersion,
1135 header,
1136 expectedId,
1137 signal,
1138 )
1139 return header
1140 }
1141
1142 /** Convert format-catalog string identities to current branded Session metadata. */
1143 private currentHeader(header: {
1144 readonly version: number
1145 readonly id: string
1146 readonly createdAt: number
1147 readonly cwd?: string
1148 readonly parentSession?: string
1149 readonly isSeeded: boolean
1150 readonly origin?: 'subagent'
1151 readonly delegationDepth: number
1152 readonly agentPreset?: string
1153 }): SessionHeader {
1154 /* v8 ignore next 3 -- readable catalog results are restored to its configured current version. */
1155 if (header.version !== SESSION_FORMAT_VERSION) {
1156 throw new Error(`format catalog returned non-current logical header v${header.version}`)
1157 }
1158 return {
1159 version: SESSION_FORMAT_VERSION,
1160 id: makeSessionId(header.id),
1161 createdAt: header.createdAt,
1162 ...(header.cwd === undefined ? {} : { cwd: header.cwd }),
1163 ...(header.parentSession === undefined
1164 ? {}
1165 : { parentSession: makeSessionId(header.parentSession) }),
1166 isSeeded: header.isSeeded,
1167 ...(header.origin === undefined ? {} : { origin: header.origin }),
1168 delegationDepth: header.delegationDepth,
1169 ...(header.agentPreset === undefined ? {} : { agentPreset: header.agentPreset }),
1170 }
1171 }
1172
1173 // --- materialization / append / repair (file mechanics) ---
1174
1175 /** Atomically write the header line + first batch (temp-write, fsync, publish). */
1176 private async materialize(
1177 meta: SessionHeader,
1178 inheritedEventCount: SessionLogOffsetType,
1179 events: readonly SessionEvent[],
1180 ): Promise<void> {
1181 const project = projectDir(this.root, meta.cwd)
1182 const dir = sessionDir(this.root, meta.cwd, meta.id)
1183 const finalPath = logPath(this.root, meta.cwd, meta.id, this.compression)
1184 await this.rejectOppositeArtifact(meta.cwd, meta.id)
1185 const content = await this.encodeMaterialization(meta, inheritedEventCount, events)
1186 /* v8 ignore next -- native Windows coverage exercises this platform dispatch; Linux covers the POSIX peer */
1187 if (process.platform === 'win32') {
1188 await this.materializeWin32(project, dir, finalPath, meta.id, content)
1189 } else {
1190 await this.materializePosix(project, dir, finalPath, meta.id, content)
1191 }
1192 }
1193
1194 /* v8 ignore start -- Windows uses the Win32 durable-publish path; POSIX coverage exercises this peer. */
1195 private async materializePosix(
1196 project: string,
1197 dir: string,
1198 finalPath: string,
1199 id: SessionId,
1200 content: Buffer | string,
1201 ): Promise<void> {
1202 await mkdir(this.root, { recursive: true, mode: 0o700 })
1203 await this.syncDirPosix(dirname(this.root))
1204 await mkdir(project, { recursive: true, mode: 0o700 })
1205 await this.syncDirPosix(this.root)
1206 await mkdir(dir, { recursive: true, mode: 0o700 })
1207 await this.syncDirPosix(project)
1208 await this.rejectExistingLog(finalPath, id)
1209 const tmp = await this.writeSyncedTempFile(finalPath, content)
1210 // Publish via link()+unlink(), NOT rename(): link fails with EEXIST if the
1211 // final path already exists, so two processes materializing the same id
1212 // concurrently cannot clobber each other. rename() would silently overwrite.
1213 let linked = false
1214 try {
1215 await link(tmp, finalPath)
1216 linked = true
1217 } finally {
1218 // Remove an unpublished temp on failure. After publication, defer cleanup
1219 // until the directory entry is durable so cleanup cannot reject a live log.
1220 /* v8 ignore next -- link failure is the TOCTOU/IO race guarded above; not reachable in test */
1221 if (!linked) await rm(tmp, { force: true })
1222 }
1223 // link() succeeded — the log is published. fsync the directory so the new
1224 // entry survives a power loss: the new link is not crash-durable until the
1225 // parent directory's metadata is synced.
1226 await this.syncDirPosix(dir)
1227 // Best-effort temp cleanup: the log is already published and durable, so a
1228 // failure to remove the redundant temp hard link must NOT reject the
1229 // append. Swallow only the rm failure; nothing else of consequence runs here.
1230 try {
1231 await rm(tmp, { force: true })
1232 } catch {
1233 /* v8 ignore next -- redundant temp link; publish already durable, rm failure is an unreachable IO edge */
1234 }
1235 }
1236 /* v8 ignore stop */
1237
1238 /* v8 ignore start -- native Windows coverage exercises this integration path */
1239 private async materializeWin32(
1240 project: string,
1241 dir: string,
1242 finalPath: string,
1243 id: SessionId,
1244 content: Buffer | string,
1245 ): Promise<void> {
1246 await ensureDurableDirectoryWin32(this.root)
1247 await ensureDurableDirectoryWin32(project)
1248 await ensureDurableDirectoryWin32(dir)
1249 await this.rejectExistingLog(finalPath, id)
1250 const tmp = await this.writeSyncedTempFile(finalPath, content)
1251 try {
1252 await publishNewFileWin32(tmp, finalPath)
1253 } catch (error) {
1254 await rm(tmp, { force: true })
1255 throw error
1256 }
1257 }
1258 /* v8 ignore stop */
1259
1260 private async rejectExistingLog(finalPath: string, id: SessionId): Promise<void> {
1261 // Never publish over an existing committed log: materialize is the first
1262 // write of a session the backend believes is new. A file here means a
1263 // different session shares this id on disk — reject loudly. (create already
1264 // guards the create path, so this is unreachable-in-practice TOCTOU
1265 // defense.)
1266 /* v8 ignore next 3 -- create guards collisions before materialize; this is a TOCTOU backstop */
1267 if (await this.resolveGenerationInDirectory(dirname(finalPath)) !== undefined) {
1268 throw new Error(`refusing to materialize "${id}": a log already exists on disk (open it instead)`)
1269 }
1270 }
1271
1272 private async writeSyncedTempFile(finalPath: string, content: Buffer | string): Promise<string> {
1273 const tmp = `${finalPath}.${randomBytes(6).toString('hex')}.tmp`
1274 const handle = await open(tmp, 'wx', 0o600)
1275 try {
1276 await handle.writeFile(content)
1277 await handle.sync()
1278 } finally {
1279 await handle.close()
1280 }
1281 return tmp
1282 }
1283
1284 /** Encode the header and first batch without combining their frame boundaries. */
1285 private async encodeMaterialization(
1286 meta: SessionHeader,
1287 inheritedEventCount: SessionLogOffsetType,
1288 events: readonly SessionEvent[],
1289 ): Promise<Buffer | string> {
1290 const header = JSON.stringify(toHeaderLine(meta, meta.isSeeded ? inheritedEventCount : undefined)) + '\n'
1291 if (events.length === 0) {
1292 return this.compression === 'none' ? header : compressZstdFrame(header)
1293 }
1294 const body = eventLines(events) + '\n'
1295 if (this.compression === 'none') return header + body
1296 const headerFrame = await compressZstdFrame(header)
1297 const eventFrame = await compressZstdFrame(body)
1298 return Buffer.concat([headerFrame, eventFrame])
1299 }
1300
1301 /** Encode one durable append batch in the configured physical representation. */
1302 private async encodeEventBatch(events: readonly SessionEvent[]): Promise<Buffer | string> {
1303 const body = eventLines(events) + '\n'
1304 return this.compression === 'zstd' ? compressZstdFrame(body) : body
1305 }
1306
1307 /** fsync a POSIX directory so a just-created/renamed entry is crash-durable. */
1308 /* v8 ignore start -- Windows uses write-through namespace operations; POSIX coverage exercises directory fsync. */
1309 private async syncDirPosix(dir: string): Promise<void> {
1310 const handle = await open(dir, 'r')
1311 try {
1312 await handle.sync()
1313 } finally {
1314 await handle.close()
1315 }
1316 }
1317 /* v8 ignore stop */
1318
1319 /**
1320 * Append and fsync event lines. On a partial write or sync failure, restore the
1321 * previous size before rethrowing because the unchanged cursor will retry the
1322 * batch; leaving partial bytes would create duplicate sequence numbers.
1323 */
1324 private async appendLines(meta: SessionHeader, events: readonly SessionEvent[]): Promise<void> {
1325 const content = await this.encodeEventBatch(events)
1326 const path = logPath(this.root, meta.cwd, meta.id, this.compression)
1327 const handle = await open(path, 'a')
1328 let closed = false
1329 const closeAppendHandle = async (): Promise<void> => {
1330 if (closed) return
1331 closed = true
1332 await handle.close()
1333 }
1334
1335 try {
1336 const { size: before } = await handle.stat()
1337 try {
1338 await handle.writeFile(content)
1339 await handle.sync()
1340 } catch (error) {
1341 try {
1342 await closeAppendHandle()
1343 await this.rollbackAppend(path, before)
1344 } catch (rollbackError) {
1345 throw new AggregateError([error, rollbackError], `failed to roll back append to "${path}"`)
1346 }
1347 throw error
1348 }
1349 } finally {
1350 await closeAppendHandle()
1351 }
1352 }
1353
1354 private async rollbackAppend(path: string, size: number): Promise<void> {
1355 const handle = await open(path, 'r+')
1356 try {
1357 await handle.truncate(size)
1358 await handle.sync()
1359 } finally {
1360 await handle.close()
1361 }
1362 }
1363
1364 /** Truncate the log file to `offset` bytes and fsync (discard the crash tail). */
1365 private async repair(meta: SessionHeader, offset: number): Promise<void> {
1366 const path = logPath(this.root, meta.cwd, meta.id, this.compression)
1367 await truncate(path, offset)
1368 const handle = await open(path, 'r+')
1369 try {
1370 await handle.sync()
1371 } finally {
1372 await handle.close()
1373 }
1374 }
1375
1376 // --- discovery helpers ---
1377
1378 /**
1379 * Read the first newline-terminated line of a file without loading the whole
1380 * file. Returns undefined if the file is empty or has no complete first line.
1381 * Reads in bounded chunks so a huge log costs only the header read.
1382 */
1383 private async readFirstLine(path: string, signal?: AbortSignal): Promise<string | undefined> {
1384 signal?.throwIfAborted()
1385 const handle = await open(path, 'r')
1386 try {
1387 signal?.throwIfAborted()
1388 const chunks: Buffer[] = []
1389 const buf = Buffer.alloc(8192)
1390 for (;;) {
1391 signal?.throwIfAborted()
1392 const { bytesRead } = await handle.read(buf, 0, buf.length, null)
1393 signal?.throwIfAborted()
1394 if (bytesRead === 0) return undefined // EOF with no newline → no complete line
1395 const slice = buf.subarray(0, bytesRead)
1396 const nl = slice.indexOf(0x0a)
1397 if (nl !== -1) {
1398 chunks.push(slice.subarray(0, nl))
1399 signal?.throwIfAborted()
1400 return Buffer.concat(chunks).toString('utf8')
1401 }
1402 chunks.push(Buffer.from(slice))
1403 }
1404 } finally {
1405 await handle.close()
1406 }
1407 }
1408
1409 /** Read only the header frame; compression failures reject as corruption, while I/O and cancellation propagate. */
1410 private async readFirstZstdLine(path: string, signal?: AbortSignal): Promise<string | undefined> {
1411 signal?.throwIfAborted()
1412 const handle = await open(path, 'r')
1413 try {
1414 signal?.throwIfAborted()
1415 let content = Buffer.alloc(0)
1416 const chunk = Buffer.alloc(8192)
1417 for (;;) {
1418 signal?.throwIfAborted()
1419 const { bytesRead } = await handle.read(chunk, 0, chunk.length, null)
1420 signal?.throwIfAborted()
1421 if (bytesRead === 0) return undefined
1422 signal?.throwIfAborted()
1423 content = Buffer.concat([content, chunk.subarray(0, bytesRead)])
1424 signal?.throwIfAborted()
1425 try {
1426 const first = scanZstdFrames(content, 1).frames[0]
1427 if (first === undefined) continue
1428 const plaintext = await decompressZstdFrame(content.subarray(first.start, first.end))
1429 signal?.throwIfAborted()
1430 assertZstdHeaderFrame(plaintext)
1431 return plaintext.subarray(0, -1).toString('utf8')
1432 } catch (error) {
1433 /* v8 ignore next -- decoder failure plus concurrent abort is timing-dependent */
1434 if (signal?.aborted) signal.throwIfAborted()
1435 throw new SessionPersistenceCorruptionError(
1436 `corrupt Zstandard session log: header frame failed validation: ${String(error)} (raw log: ${path})`,
1437 { cause: error },
1438 )
1439 }
1440 }
1441 } finally {
1442 await handle.close()
1443 }
1444 }
1445
1446 /** Select the numerically highest canonical generation in one Session directory. */
1447 private async resolveGenerationInDirectory(
1448 dir: string,
1449 signal?: AbortSignal,
1450 ): Promise<ResolvedJsonlGeneration | undefined> {
1451 signal?.throwIfAborted()
1452 let entries: Dirent[]
1453 try {
1454 entries = await readdir(dir, { withFileTypes: true })
1455 } catch (error: unknown) {
1456 if (isENOENT(error)) return undefined
1457 throw error
1458 }
1459 signal?.throwIfAborted()
1460 const generations: Array<{ readonly path: string; readonly version: number }> = []
1461 const opposite: string[] = []
1462 for (const entry of entries) {
1463 const version = parseGenerationLogFilename(entry.name, this.compression)
1464 if (version !== undefined) {
1465 generations.push({ path: join(dir, entry.name), version })
1466 continue
1467 }
1468 if (parseGenerationLogFilename(entry.name, this.oppositeCompression()) !== undefined) {
1469 opposite.push(join(dir, entry.name))
1470 }
1471 }
1472 if (opposite.length > 0) throw this.encodingMismatch(opposite[0] as string)
1473 const latest = generations.sort((left, right) => right.version - left.version)[0]
1474 if (latest === undefined) return undefined
1475 return {
1476 sourcePath: latest.path,
1477 sourceVersion: latest.version,
1478 currentPath: join(
1479 dir,
1480 generationLogFilename(sessionFormatCatalog.currentVersion, this.compression),
1481 ),
1482 }
1483 }
1484
1485 /** Find the unique authoritative generation for an id across project directories. */
1486 private async findLog(id: SessionId, signal?: AbortSignal): Promise<ResolvedJsonlGeneration | undefined> {
1487 const matches: ResolvedJsonlGeneration[] = []
1488 for (const project of await this.listProjectDirs(signal)) {
1489 signal?.throwIfAborted()
1490 await this.rejectLegacyFlatArtifact(project, id, signal)
1491 signal?.throwIfAborted()
1492 const dir = join(project, encodeSegment(id))
1493 const selected = await this.resolveGenerationInDirectory(dir, signal)
1494 if (selected !== undefined) matches.push(selected)
1495 }
1496 if (matches.length > 1) {
1497 throw new Error(`duplicate JSONL session id "${id}" appears in multiple project directories`)
1498 }
1499 signal?.throwIfAborted()
1500 return matches[0]
1501 }
1502
1503 /** Require an existing configured root to be a readable directory. */
1504 private assertUsableRoot(): void {
1505 try {
1506 readdirSync(this.root)
1507 } catch (error) {
1508 if (isENOENT(error)) return
1509 throw error
1510 }
1511 }
1512
1513 /** Reject metadata that does not identify the selected physical log. */
1514 private async assertStoredIdentity(
1515 path: string,
1516 storedVersion: number,
1517 meta: SessionHeader,
1518 expectedId?: SessionId,
1519 signal?: AbortSignal,
1520 ): Promise<void> {
1521 signal?.throwIfAborted()
1522 if (expectedId !== undefined && meta.id !== expectedId) {
1523 throw new Error(`corrupt session log "${path}": requested id "${expectedId}" does not match header id "${meta.id}"`)
1524 }
1525 let expectedPath: string
1526 try {
1527 expectedPath = generationLogPath(
1528 this.root,
1529 meta.cwd,
1530 meta.id,
1531 storedVersion,
1532 this.compression,
1533 )
1534 } catch (error) {
1535 throw new Error(`corrupt session log "${path}": header id cannot name a storage path`, { cause: error })
1536 }
1537 if (path !== expectedPath && !await this.sameFile(path, expectedPath, signal)) {
1538 throw new Error(`corrupt session log "${path}": header id "${meta.id}" and cwd identify "${expectedPath}"`)
1539 }
1540 signal?.throwIfAborted()
1541 }
1542
1543 /** Validate a supported historical header against the selected source path. */
1544 private validateSourceIdentity(
1545 selected: ResolvedJsonlGeneration,
1546 headerValue: Readonly<Record<string, unknown>>,
1547 expectedId: SessionId,
1548 signal?: AbortSignal,
1549 ): void | Promise<void> {
1550 const result = sessionFormatCatalog.readHeader(headerValue)
1551 if (result.status !== 'current' && result.status !== 'migration-required') return
1552 return this.assertStoredIdentity(
1553 selected.sourcePath,
1554 selected.sourceVersion,
1555 this.currentHeader(result.header),
1556 expectedId,
1557 signal,
1558 )
1559 }
1560
1561 /**
1562 * Whether two path spellings resolve to the same physical file. This admits
1563 * case aliases on case-insensitive filesystems without weakening identity
1564 * checks on case-sensitive stores.
1565 */
1566 private async sameFile(path: string, expectedPath: string, signal?: AbortSignal): Promise<boolean> {
1567 signal?.throwIfAborted()
1568 try {
1569 const [actual, expected] = await Promise.all([realpath(path), realpath(expectedPath)])
1570 signal?.throwIfAborted()
1571 return actual === expected
1572 } catch (error) {
1573 signal?.throwIfAborted()
1574 /* v8 ignore else -- non-ENOENT realpath failures require an external permission or I/O fault */
1575 if (isENOENT(error)) return false
1576 /* v8 ignore next -- non-ENOENT realpath failures are external I/O faults, propagated unchanged */
1577 throw error
1578 }
1579 }
1580
1581 /** The human-readable project directories under the configured root. */
1582 private async listProjectDirs(signal?: AbortSignal): Promise<string[]> {
1583 try {
1584 signal?.throwIfAborted()
1585 const entries = await readdir(this.root, { withFileTypes: true })
1586 signal?.throwIfAborted()
1587 return entries.filter(e => e.isDirectory()).map(e => join(this.root, e.name))
1588 } catch (error) {
1589 // Only an absent root means no sessions; rethrow every other I/O failure.
1590 if (isENOENT(error)) return []
1591 throw error
1592 }
1593 }
1594
1595 /** List session-owned directories and reject the obsolete flat-file layout. */
1596 private async listSessionDirs(project: string, signal?: AbortSignal): Promise<string[]> {
1597 signal?.throwIfAborted()
1598 const entries = await readdir(project, { withFileTypes: true })
1599 signal?.throwIfAborted()
1600 const legacy = entries.find(entry =>
1601 entry.isFile() && (entry.name.endsWith('.jsonl') || entry.name.endsWith('.jsonl.zstd')))
1602 if (legacy !== undefined) throw this.legacyLayout(join(project, legacy.name))
1603 return entries.filter(entry => entry.isDirectory()).map(entry => join(project, entry.name))
1604 }
1605
1606 /** Reject a root that already belongs to the other physical encoding. */
1607 private ensureRootEncoding(): Promise<void> {
1608 this.rootEncodingCheck ??= this.checkRootEncoding()
1609 return this.rootEncodingCheck
1610 }
1611
1612 private async checkRootEncoding(): Promise<void> {
1613 for (const project of await this.listProjectDirs()) {
1614 for (const dir of await this.listSessionDirs(project)) {
1615 const incompatible = await this.findOppositeGenerationInDirectory(dir)
1616 if (incompatible !== undefined) throw this.encodingMismatch(incompatible)
1617 }
1618 }
1619 }
1620
1621 private async rejectLegacyFlatArtifact(
1622 project: string,
1623 id: SessionId,
1624 signal?: AbortSignal,
1625 ): Promise<void> {
1626 signal?.throwIfAborted()
1627 const encoded = encodeSegment(id)
1628 for (const compression of ['zstd', 'none'] as const) {
1629 const path = join(project, encoded + logSuffix(compression))
1630 const artifactExists = await this.exists(path)
1631 signal?.throwIfAborted()
1632 if (artifactExists) throw this.legacyLayout(path)
1633 }
1634 }
1635
1636 private async rejectOppositeArtifact(cwd: string | undefined, id: SessionId): Promise<void> {
1637 const path = await this.findOppositeGenerationInDirectory(sessionDir(this.root, cwd, id))
1638 if (path !== undefined) throw this.encodingMismatch(path)
1639 }
1640
1641 /** Return the highest canonical generation encoded with the other configured suffix. */
1642 private async findOppositeGenerationInDirectory(dir: string): Promise<string | undefined> {
1643 let entries: Dirent[]
1644 try {
1645 entries = await readdir(dir, { withFileTypes: true })
1646 } catch (error: unknown) {
1647 if (isENOENT(error)) return undefined
1648 throw error
1649 }
1650 const generations: Array<{ readonly name: string; readonly version: number }> = []
1651 for (const entry of entries) {
1652 const version = parseGenerationLogFilename(entry.name, this.oppositeCompression())
1653 if (version !== undefined) generations.push({ name: entry.name, version })
1654 }
1655 const latest = generations.sort((left, right) => right.version - left.version)[0]
1656 return latest === undefined ? undefined : join(dir, latest.name)
1657 }
1658
1659 private oppositeCompression(): JsonlCompression {
1660 return this.compression === 'zstd' ? 'none' : 'zstd'
1661 }
1662
1663 private encodingMismatch(path: string): Error {
1664 return new Error(
1665 `session artifact ${JSON.stringify(path)} uses ${logSuffix(this.oppositeCompression())}, `
1666 + `but this backend is configured for compression ${JSON.stringify(this.compression)}; `
1667 + 'use a separate root or select the matching compression mode',
1668 )
1669 }
1670
1671 private legacyLayout(path: string): Error {
1672 return new Error(
1673 `session artifact ${JSON.stringify(path)} uses the unsupported flat-file layout; `
1674 + 'use a separate root or move it into a project/session directory before loading',
1675 )
1676 }
1677
1678 private async exists(path: string): Promise<boolean> {
1679 try {
1680 const handle = await open(path, 'r')
1681 await handle.close()
1682 return true
1683 } catch (error) {
1684 // Only ENOENT means absent. A permission/I/O error must surface rather
1685 // than letting load or collision checks proceed under false absence.
1686 /* v8 ignore else -- Windows reports file-valued parents as ENOENT; POSIX covers direct ENOTDIR. */
1687 if (isENOENT(error)) {
1688 // Windows reports ENOENT, not ENOTDIR, for `regular-file/child`, so it
1689 // alone verifies the immediate parent to keep a blocked session
1690 // directory a storage fault. POSIX open already reported ENOTDIR before
1691 // this point, where the extra stat would only cost a syscall per probe.
1692 /* v8 ignore next -- native Windows coverage exercises this platform dispatch; POSIX reports ENOTDIR from open */
1693 if (process.platform === 'win32') await this.assertLogParentAllowsAbsence(path)
1694 return false
1695 }
1696 /* v8 ignore next -- Windows repairs ENOTDIR from ENOENT above; POSIX covers direct ENOTDIR. */
1697 throw error
1698 }
1699 }
1700
1701 /* v8 ignore start -- native Windows coverage exercises this repair; POSIX open reports ENOTDIR before this point. */
1702 private async assertLogParentAllowsAbsence(path: string): Promise<void> {
1703 try {
1704 const parent = dirname(path)
1705 const info = await stat(parent)
1706 if (info.isDirectory()) return
1707 const error = new Error(`ENOTDIR: parent path exists but is not a directory: ${parent}`) as NodeJS.ErrnoException
1708 error.code = 'ENOTDIR'
1709 error.path = parent
1710 throw error
1711 } catch (error) {
1712 if (isENOENT(error)) return
1713 throw error
1714 }
1715 }
1716 /* v8 ignore stop */
1717}
1718
1719/**
1720 * One open channel onto a JSONL-stored session: the shared storage-handle
1721 * scaffolding over this backend's file primitives. Reads re-scan the artifact
1722 * under the stable-read loop.
1723 */
1724
1725export default JsonlSessionPersistence