1
/**2
* JSONL durable session-persistence backend. It stores a header and contiguous3
* events in immutable generation files under one directory per session and serves the handle-based4
* `SessionPersistence` API: `create`/`open` return per-session handles, and5
* every read validates the same fail-closed storage contract.6
* @module @deepseek-ai/dsh-session-persistence-jsonl7
*/9
import { Context } from '@deepseek-ai/cordis'10
import z from '@deepseek-ai/schemastery'11
import {12
createSessionFormatCatalogWithChildren,13
SessionFormatUnsupportedMigrationError,14
sessionFormatCatalog,15
} from '@deepseek-ai/dsh-session-format-catalog'16
import { readdirSync, type Dirent } from 'node:fs'17
import { open, mkdir, readdir, realpath, link, rm, stat, truncate } from 'node:fs/promises'18
import { dirname, join, resolve } from 'node:path'19
import { performance } from 'node:perf_hooks'20
import { scheduler } from 'node:timers/promises'21
import { createHash, randomBytes } from 'node:crypto'22
import {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'34
import { JsonlBackendTracker, JsonlSessionHandle, type StorageHandleState } from './storage.ts'35
import { SessionWriteLease } from './lease.ts'36
import { SESSION_FORMAT_VERSION, SessionId as makeSessionId, SessionLogOffset } from '@deepseek-ai/dsh-session'37
import type { SessionEvent, SessionId, SessionHeader, SessionLogOffset as SessionLogOffsetType } from '@deepseek-ai/dsh-session'38
import {39
assertNoRetiredHeaderFields, encodeSegment, eventLines, generationLogFilename, generationLogPath, logPath, logSuffix,40
parseGenerationLogFilename, projectDir, scanLog, sessionDir, SessionLogScanner, toHeaderLine,41
type JsonlCompression,42
} from './format.ts'43
import {44
compressZstdFrame, createZstdFrameDecoder, decompressZstdFrame, decompressZstdPrefix, scanZstdFrames,45
} from './zstd.ts'46
import { ensureDurableDirectoryWin32, publishNewFileWin32 } from './win32.ts'47
import { verifyCurrentGenerationInWorker } from './migration-verifier.ts'48
import { prepareCatalogFacts } from './catalog-migration.ts'49
import {50
JsonlGenerationSourceChangedError,51
JsonlGenerationUnsupportedMigrationError,52
prepareJsonlMigration,53
readStableJsonlFile,54
type JsonlGenerationFormatAdapter,55
type JsonlPhysicalIdentity,56
type PreparedJsonlMigration,57
} from './generation.ts'59
export type { JsonlCompression } from './format.ts'61
/**62
* Internal handoff-reuse policy, not deployment configuration: a cold63
* observation and the resume that immediately follows it reuse one parsed64
* log, so the memo only needs the sessions in flight between those steps.65
*/66
const COLD_LOG_MEMO_MAX_ENTRIES = 268
const DEFAULT_COMPRESSION: JsonlCompression = 'zstd'69
/**70
* Internal scheduling constant, not deployment configuration: balance71
* frame-boundary event-loop yields against `setImmediate` overhead. One frame72
* remains an indivisible synchronous decode.73
*/74
const ZSTD_DECODE_YIELD_INTERVAL_MS = 50076
/** Assert that the independently decodable first frame contains only the header record. */77
function 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
}83
/** Loader schema for the JSONL artifact's physical encoding. */84
export const JsonlCompressionSchema: z<JsonlCompression> = z.union([85
z.const('zstd'),86
z.const('none'),87
]).default(DEFAULT_COMPRESSION)89
/** Plugin config for the JSONL backend's root and physical encoding. */90
export interface Config {91
/**92
* Root directory for all session files. Required (no default): a default of93
* `process.cwd()` would scatter session files as the process's cwd changes94
* (bash calls, subprocesses). Sessions group under human-readable project95
* directories, then per-session directories. An existing root must be a96
* readable directory; an absent root is created on first materialization.97
*/98
root: string99
/** Physical encoding; defaults to checksummed Zstandard frames. */100
compression?: JsonlCompression101
}103
/** One stored event graph whose producer has established immutable sharing. */104
interface FrozenStoredEvents extends SessionHandleReadResult {105
readonly eventState: 'shared-frozen'106
}108
/** State shared by prepared historical and published current logs. */109
interface StoredLogBase extends FrozenStoredEvents {110
readonly meta: SessionHeader111
readonly tornTruncateTo: number | undefined112
/** 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: SessionLogOffsetType116
readonly revision: PersistenceRevision117
}119
/** A decoded current generation that is already durable. */120
interface CurrentStoredLog extends StoredLogBase {121
readonly status: 'current'122
}124
/** A migrated historical generation retained until an explicit write open publishes it. */125
interface PreparedStoredLog extends StoredLogBase {126
readonly status: 'prepared'127
readonly validateRelatedSources: () => Promise<void>128
readonly publication: {129
readonly source: ResolvedJsonlGeneration130
readonly value: PreparedJsonlMigration131
}132
}134
/** A validated logical log, either durable current state or prepared historical state. */135
type StoredLog = CurrentStoredLog | PreparedStoredLog137
/** Deep-freeze acyclic stored JSON; its arrays contain only indexed JSON values. */138
function 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-assertion143
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
}159
/** Establish immutable sharing for one decoded event graph and report that state. */160
function freezeStoredEvents(events: SessionEvent[]): FrozenStoredEvents {161
for (const event of events) freezeStoredEvent(event)162
Object.freeze(events)163
return { eventState: 'shared-frozen', events }164
}166
/** One authoritative immutable generation selected from a Session directory. */167
interface ResolvedJsonlGeneration {168
readonly sourcePath: string169
readonly sourceVersion: number170
readonly currentPath: string171
}173
/** One backend-owned historical preparation shared by its current callers. */174
interface MigrationPreparation {175
readonly sourcePath: string176
readonly sourceRevision: PersistenceRevision177
readonly controller: AbortController178
readonly promise: Promise<PreparedStoredLog>179
settled: boolean180
waiters: number181
}183
/** Build the stat-derived best-effort change token shared by full and lightweight reads. */184
function 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
}194
/** Whether a filesystem error means absence; every non-ENOENT failure must surface. */195
function isENOENT(error: unknown): boolean {196
return (error as NodeJS.ErrnoException | null)?.code === 'ENOENT'197
}199
/** Whether a filesystem-owned failure should retain its original errno and path. */200
function isErrnoException(error: unknown): error is NodeJS.ErrnoException {201
return typeof (error as NodeJS.ErrnoException | null)?.code === 'string'202
}204
/** Preserve an Error abort reason and normalize hostile non-Error reasons. */205
function abortError(signal: AbortSignal): Error {206
return signal.reason instanceof Error207
? signal.reason208
: new Error('session migration preparation aborted', { cause: signal.reason })209
}211
/** Let one caller stop waiting without transferring cancellation ownership to shared work. */212
function waitWithAbort<T>(operation: Promise<T>, signal?: AbortSignal): Promise<T> {213
if (signal === undefined) return operation214
/* 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
}239
/**240
* The JSONL persistence backend. Load as a plugin; it registers as241
* `ctx.sessionPersistence`. Sessions materialize lazily: a created session is242
* visible to this process immediately, reaches disk on its first append or243
* flush, and never existed if the process crashes before that.244
*/245
class JsonlSessionPersistence extends SessionPersistence {246
static Config: z<Config> = z.object({247
root: z.string().required(),248
compression: JsonlCompressionSchema,249
})251
/** Backend label for diagnostics and effects; shadows `Service.name` without changing the service key. */252
override readonly name = 'session-persistence-jsonl'254
private root: string255
private compression: JsonlCompression256
private rootEncodingCheck: Promise<void> | undefined257
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 and261
* guarded by the stat-derived revision, so an immediate cold-read handoff262
* (observation then resume) parses the artifact once. Every local mutation263
* for an id invalidates its entry; a foreign write misses through the264
* 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>()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_COMPRESSION282
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
}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
}303
// --- SessionPersistence service API ---305
/**306
* Create a new stored session and take its write ownership. The session is307
* visible to this process immediately; the physical artifact appears on the308
* first append or flush.309
* @param header - the immutable header to store; must be losslessly310
* 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 header318
// 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 for328
// another process to contend over, so the handle acquires the lock right329
// before its first log bytes publish (ensureLease); an unmaterialized330
// 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
}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: StoredLog352
try {353
stored = await this.requireStoredLog(id, options?.signal)354
} catch (error: unknown) {355
if (!(error instanceof JsonlGenerationSourceChangedError)) throw error356
stored = await this.requireStoredLog(id, options?.signal)357
}358
let state: StorageHandleState359
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 erases376
// it), so the claim below rejects that case as already owned.377
this.tracker.claimWrite(id)378
let lease: SessionWriteLease | undefined379
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: CurrentStoredLog386
if (prepared.status === 'prepared') {387
stored = await this.publishStoredMigration(id, prepared)388
} else {389
stored = prepared390
}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 release402
// fares, and keep the original diagnostic: a release failure joins it403
// 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 | undefined407
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 failure418
}419
}421
/**422
* Flush every active write handle in one durability barrier; see the seam423
* 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
}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), or435
* `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 undefined450
const header = await this.readGenerationHeader(selected, id, options?.signal)451
if (header === undefined) return undefined452
try {453
const identity = await stat(selected.sourcePath, { bigint: true })454
options?.signal?.throwIfAborted()455
return {456
header,457
revision: selected.sourceVersion < SESSION_FORMAT_VERSION458
? 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 undefined465
throw error466
}467
}469
/**470
* List every stored session visible to this process: materialized artifacts471
* 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?.signal477
const snapshots: SessionPersistenceSnapshot[] = []478
const listed = new Set<SessionId>()479
// Snapshot pending entries BEFORE scanning storage: a session whose first480
// append lands mid-scan is then still in this snapshot (its artifact may481
// 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) : undefined486
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_VERSION495
? 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 error502
}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 snapshots509
}511
// --- handle-facing storage internals (package-private via the handle class below) ---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 === undefined522
|| preparation.sourcePath !== selected.sourcePath523
|| 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 = preparation536
const release = (): void => {537
created.settled = true538
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 memoized568
}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
}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 memoized598
}599
return this.prepareStoredMigration(id, selected, signal)600
}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 += 1609
try {610
return await waitWithAbort(preparation.promise, signal)611
} finally {612
preparation.waiters -= 1613
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
}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 stored689
}691
/** Publish a prepared historical log before granting write access. */692
private async publishStoredMigration(id: SessionId, stored: PreparedStoredLog): Promise<CurrentStoredLog> {693
const migration = stored.publication694
let identity: JsonlPhysicalIdentity695
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 published714
}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 error729
if (error instanceof SessionFormatUnsupportedError730
|| error instanceof SessionPersistenceCorruptionError731
|| isErrnoException(error)732
|| error instanceof DOMException && error.name === 'AbortError') return error733
return new SessionPersistenceCorruptionError(734
`session "${id}": stored log is corrupt: ${String(error)} (raw log: ${selected.sourcePath})`,735
{ cause: error },736
)737
}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 memoized754
}755
const { bytes, identity } = await readStableJsonlFile(path, signal)756
return this.decodeStoredLog(path, expectedId, bytes, fileRevision(identity), signal)757
}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: SessionHeader769
inheritedEventCount: SessionLogOffsetType770
events: SessionEvent[]771
tornTruncateTo: number | undefined772
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 complete787
// 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 the794
// artifact this read actually refused; every other parse failure is795
// committed bytes the decoder cannot interpret — damage, classified for796
// 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 } = parsed809
const stored: CurrentStoredLog = {810
status: 'current',811
...rest,812
...freezeStoredEvents(events),813
revision,814
}815
this.memoizeStoredLog(expectedId, stored)816
return stored817
}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) break825
this.coldLogMemo.delete(oldest)826
}827
}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 undefined840
if (selected.sourceVersion === SESSION_FORMAT_VERSION) return selected.sourcePath841
if (selected.sourceVersion < SESSION_FORMAT_VERSION) return undefined842
const reason = sessionFormatVersionRefusal(id, selected.sourceVersion)843
throw new SessionFormatUnsupportedError(844
`${reason} (raw log: ${selected.sourcePath})`,845
{ kind: 'jsonl', path: selected.sourcePath },846
)847
}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
}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
}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
}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
}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
}913
/**914
* Acquire the session directory's kernel write lock; the kernel holds it915
* 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
}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 publishes933
// 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
}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: SessionHeader944
inheritedEventCount: SessionLogOffsetType945
events: SessionEvent[]946
tornTruncateTo: number | undefined947
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')954
const decoder = createZstdFrameDecoder()955
let yieldDeadline = performance.now() + ZSTD_DECODE_YIELD_INTERVAL_MS956
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)966
let remainingFrames = frames.length - 1967
for (const plaintext of decodedFrames) {968
signal?.throwIfAborted()969
scanner.write(plaintext)970
remainingFrames -= 1971
if (remainingFrames > 0 && performance.now() >= yieldDeadline) {972
await scheduler.yield()973
signal?.throwIfAborted()974
yieldDeadline = performance.now() + ZSTD_DECODE_YIELD_INTERVAL_MS975
}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 records993
// already flushed into it are real emitted events: recover them, and let994
// 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 decoder1003
// 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 error1019
} finally {1020
decoder.close()1021
}1022
}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 sources1035
}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: string1044
try {1045
revision = fileRevision(await stat(path, { bigint: true }))1046
} catch (error: unknown) {1047
if (!isENOENT(error)) throw error1048
revision = 'missing'1049
}1050
hash.update(JSON.stringify([path, revision]))1051
}1052
signal?.throwIfAborted()1053
return hash.digest('hex')1054
}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 | undefined1067
try {1068
header = await this.readGenerationHeader(selected, undefined, signal)1069
} catch (error: unknown) {1070
if (error instanceof SessionFormatUnsupportedError || error instanceof SessionPersistenceCorruptionError) continue1071
throw error1072
}1073
if (header === undefined) {1074
continue1075
}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 artifacts1084
}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 | undefined1093
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 undefined1100
throw error1101
}1102
signal?.throwIfAborted()1103
if (first === undefined) return undefined1104
let value: unknown1105
try {1106
value = JSON.parse(first)1107
} catch {1108
return undefined1109
}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.reason1121
/* 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 undefined1131
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 header1140
}1142
/** Convert format-catalog string identities to current branded Session metadata. */1143
private currentHeader(header: {1144
readonly version: number1145
readonly id: string1146
readonly createdAt: number1147
readonly cwd?: string1148
readonly parentSession?: string1149
readonly isSeeded: boolean1150
readonly origin?: 'subagent'1151
readonly delegationDepth: number1152
readonly agentPreset?: string1153
}): 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 === undefined1164
? {}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
}1173
// --- materialization / append / repair (file mechanics) ---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
}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 the1211
// final path already exists, so two processes materializing the same id1212
// concurrently cannot clobber each other. rename() would silently overwrite.1213
let linked = false1214
try {1215
await link(tmp, finalPath)1216
linked = true1217
} finally {1218
// Remove an unpublished temp on failure. After publication, defer cleanup1219
// 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 new1224
// entry survives a power loss: the new link is not crash-durable until the1225
// parent directory's metadata is synced.1226
await this.syncDirPosix(dir)1227
// Best-effort temp cleanup: the log is already published and durable, so a1228
// failure to remove the redundant temp hard link must NOT reject the1229
// 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 */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 error1256
}1257
}1258
/* v8 ignore stop */1260
private async rejectExistingLog(finalPath: string, id: SessionId): Promise<void> {1261
// Never publish over an existing committed log: materialize is the first1262
// write of a session the backend believes is new. A file here means a1263
// different session shares this id on disk — reject loudly. (create already1264
// guards the create path, so this is unreachable-in-practice TOCTOU1265
// 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
}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 tmp1282
}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 + body1296
const headerFrame = await compressZstdFrame(header)1297
const eventFrame = await compressZstdFrame(body)1298
return Buffer.concat([headerFrame, eventFrame])1299
}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) : body1305
}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 */1319
/**1320
* Append and fsync event lines. On a partial write or sync failure, restore the1321
* previous size before rethrowing because the unchanged cursor will retry the1322
* 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 = false1329
const closeAppendHandle = async (): Promise<void> => {1330
if (closed) return1331
closed = true1332
await handle.close()1333
}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 error1348
}1349
} finally {1350
await closeAppendHandle()1351
}1352
}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
}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
}1376
// --- discovery helpers ---1378
/**1379
* Read the first newline-terminated line of a file without loading the whole1380
* 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 line1395
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
}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 undefined1422
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) continue1428
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
}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 undefined1457
throw error1458
}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
continue1467
}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 undefined1475
return {1476
sourcePath: latest.path,1477
sourceVersion: latest.version,1478
currentPath: join(1479
dir,1480
generationLogFilename(sessionFormatCatalog.currentVersion, this.compression),1481
),1482
}1483
}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
}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)) return1509
throw error1510
}1511
}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: string1526
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
}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') return1552
return this.assertStoredIdentity(1553
selected.sourcePath,1554
selected.sourceVersion,1555
this.currentHeader(result.header),1556
expectedId,1557
signal,1558
)1559
}1561
/**1562
* Whether two path spellings resolve to the same physical file. This admits1563
* case aliases on case-insensitive filesystems without weakening identity1564
* 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 === expected1572
} 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 false1576
/* v8 ignore next -- non-ENOENT realpath failures are external I/O faults, propagated unchanged */1577
throw error1578
}1579
}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 error1592
}1593
}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
}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.rootEncodingCheck1610
}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
}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
}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
}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 undefined1648
throw error1649
}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
}1659
private oppositeCompression(): JsonlCompression {1660
return this.compression === 'zstd' ? 'none' : 'zstd'1661
}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
}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
}1678
private async exists(path: string): Promise<boolean> {1679
try {1680
const handle = await open(path, 'r')1681
await handle.close()1682
return true1683
} catch (error) {1684
// Only ENOENT means absent. A permission/I/O error must surface rather1685
// 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 it1689
// alone verifies the immediate parent to keep a blocked session1690
// directory a storage fault. POSIX open already reported ENOTDIR before1691
// 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 false1695
}1696
/* v8 ignore next -- Windows repairs ENOTDIR from ENOENT above; POSIX covers direct ENOTDIR. */1697
throw error1698
}1699
}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()) return1707
const error = new Error(`ENOTDIR: parent path exists but is not a directory: ${parent}`) as NodeJS.ErrnoException1708
error.code = 'ENOTDIR'1709
error.path = parent1710
throw error1711
} catch (error) {1712
if (isENOENT(error)) return1713
throw error1714
}1715
}1716
/* v8 ignore stop */1717
}1719
/**1720
* One open channel onto a JSONL-stored session: the shared storage-handle1721
* scaffolding over this backend's file primitives. Reads re-scan the artifact1722
* under the stable-read loop.1723
*/1725
export default JsonlSessionPersistence