1
/**2
* Durable whole-generation publication for JSONL Session artifacts.3
*4
* Format packages transform parsed JSON values. This module owns the physical5
* encoding, exact source identity, immutable generation files, and exclusive6
* current-generation publication for both configured JSONL suffixes.7
* @module @deepseek-ai/dsh-session-persistence-jsonl/generation8
*/10
import { currentSessionMessageProjections } from '@deepseek-ai/dsh-session-format-catalog/message-projections'11
import { createHash, randomBytes } from 'node:crypto'12
import {13
link as fsLink,14
lstat as fsLstat,15
open as fsOpen,16
readFile as fsReadFile,17
readdir as fsReaddir,18
rm as fsRm,19
stat as fsStat,20
type FileHandle,21
} from 'node:fs/promises'22
import { basename, dirname, join } from 'node:path'23
import { performance } from 'node:perf_hooks'24
import { pipeline, Readable } from 'node:stream'25
import { scheduler } from 'node:timers/promises'26
import { isDeepStrictEqual } from 'node:util'27
import { constants, createZstdCompress } from 'node:zlib'28
import { Session } from '@deepseek-ai/dsh-session'29
import type { SessionEvent } from '@deepseek-ai/dsh-session'30
import { BlockAssembler, expandAssistantStream } from '@deepseek-ai/dsh-llm'31
import { SessionFormatError } from '@deepseek-ai/dsh-session-format'32
import type {33
SessionFormatArtifact,34
SessionFormatJsonValue,35
SessionFormatRestore,36
} from '@deepseek-ai/dsh-session-format'37
import { validateStoredEvents } from '@deepseek-ai/dsh-session-persistence'38
import type { JsonlCompression } from './format.ts'39
import { generationLogFilename, logSuffix, SessionLogScanner } from './format.ts'40
import { publishNewFileWin32 } from './win32.ts'41
import {42
compressZstdFrame,43
createZstdFrameDecoder,44
decompressZstdPrefix,45
scanZstdFrames,46
} from './zstd.ts'48
/** Internal scheduling bounds: preserve old decode cadence and cap each synchronous encode slice. */49
const MIGRATION_DECODE_YIELD_INTERVAL_MS = 50050
const MIGRATION_WORK_CHUNK_BYTES = 1024 * 102451
const MIGRATION_WRITE_CHUNK_BYTES = 4 * 1024 * 102452
const ZSTD_CHECKSUM_OPTIONS = {53
chunkSize: MIGRATION_WORK_CHUNK_BYTES,54
params: { [constants.ZSTD_c_checksumFlag]: 1 },55
}57
/** Pure adapter between backend-owned JSONL framing and the format catalog. */58
export interface JsonlGenerationFormatAdapter {59
readonly currentVersion: number60
/** Create the single-pass codec and migration state for a historical header. */61
createRestore(header: Record<string, unknown>): SessionFormatRestore62
/** Encode one current header record without materializing body rows. */63
encodeHeader(header: SessionFormatArtifact['header'], inheritedEventCount: number): SessionFormatJsonValue64
/** Encode one current event record. */65
encodeEvent(event: SessionFormatArtifact['events'][number]): SessionFormatJsonValue66
/** Classify a supported-version artifact that policy refuses to migrate. */67
isUnsupportedMigrationError?(error: unknown): error is Error68
}70
/** Inputs for preparing one historical generation and publishing its current successor later. */71
export interface PrepareJsonlMigrationOptions {72
/** Revalidate related source facts before preparation returns and immediately before publication. */73
readonly validateRelatedSources?: () => Promise<void>74
/** Immutable generation selected by the backend resolver. */75
readonly sourcePath: string76
/** Version selected from the source filename and independently checked against its header. */77
readonly sourceVersion: number78
/** Canonical filename for `format.currentVersion` in the same Session directory. */79
readonly currentPath: string80
readonly compression: JsonlCompression81
readonly format: JsonlGenerationFormatAdapter82
/** Validate one selected historical header's identity before any migration write. */83
readonly validateHistoricalHeader?: (84
header: Readonly<Record<string, unknown>>,85
) => void | Promise<void>86
/** Validate the staged file in an isolated worker before publication. */87
readonly verifyCurrentFile: (88
path: string,89
compression: JsonlCompression,90
expectedId: string,91
expectedEventCount: number,92
expectedPrefix?: JsonlExpectedPrefix,93
signal?: AbortSignal,94
) => Promise<JsonlVerifiedGeneration>95
readonly signal?: AbortSignal96
}98
/** Small physical identity returned by an isolated generation verifier. */99
export interface JsonlVerifiedGeneration {100
readonly identity: JsonlPhysicalIdentity101
readonly bytes: number102
readonly digest: string103
}105
/** Physical byte prefix already proven to be a valid complete generation. */106
export interface JsonlExpectedPrefix {107
readonly bytes: number108
readonly digest: string109
}111
/** A historical source changed after its single decode and migration pass. */112
export class JsonlGenerationSourceChangedError extends Error {113
override readonly name = 'JsonlGenerationSourceChangedError'115
/** @param path - historical generation whose revision changed. */116
constructor(readonly path: string) {117
super(`historical session generation changed during migration: "${path}"`)118
}119
}121
/** Current logical state prepared independently from durable publication. */122
export interface PreparedJsonlMigration {123
readonly sourceIdentity: JsonlPhysicalIdentity124
readonly artifact: SessionFormatArtifact125
/** Encode, verify, and exclusively publish once; every call shares the same success or failure. */126
publish(): Promise<JsonlPhysicalIdentity>127
}129
/** A historical artifact is intact, but the format edge refuses its contents. */130
export class JsonlGenerationUnsupportedMigrationError extends Error {131
override readonly name = 'JsonlGenerationUnsupportedMigrationError'133
/**134
* @param fromVersion - unchanged source generation version.135
* @param reason - format-edge refusal.136
*/137
constructor(138
readonly fromVersion: number,139
readonly reason: Error,140
) {141
super(reason.message, { cause: reason })142
}143
}145
/** A current-generation filename already names different or invalid bytes. */146
export class JsonlGenerationTargetConflictError extends Error {147
override readonly name = 'JsonlGenerationTargetConflictError'149
/**150
* @param path - immutable target that prevented exclusive publication.151
* @param reason - why the existing target cannot be accepted.152
*/153
constructor(154
readonly path: string,155
readonly reason: Error,156
) {157
super(`current session generation already exists at "${path}": ${reason.message}`, { cause: reason })158
}159
}161
/** Stat identity captured together with exact generation bytes. */162
export interface JsonlPhysicalIdentity {163
readonly dev: bigint164
readonly ino: bigint165
readonly size: bigint166
readonly mtimeNs: bigint167
readonly ctimeNs: bigint168
}170
/** Exact bytes of one stable file revision together with the stat identity that proved it stable. */171
export interface StablePhysicalFile {172
readonly bytes: Buffer173
readonly identity: JsonlPhysicalIdentity174
}176
interface GenerationFileSystem {177
open(path: string, flags: string, mode?: number): Promise<FileHandle>178
readFile(path: string, signal?: AbortSignal): Promise<Buffer>179
readdir(path: string): Promise<string[]>180
stat(path: string): Promise<JsonlPhysicalIdentity>181
lstat(path: string): Promise<{ isFile(): boolean; isSymbolicLink(): boolean }>182
link(existingPath: string, newPath: string): Promise<void>183
rm(path: string): Promise<void>184
}186
type GenerationBarrierPhase =187
| 'before-source-check'188
| 'after-publication'190
interface JsonlGenerationInternals {191
readonly fs: GenerationFileSystem192
readonly randomToken: () => string193
readonly platform: NodeJS.Platform194
readonly publishNewWin32: typeof publishNewFileWin32195
readonly barrier: (phase: GenerationBarrierPhase, attempt: number) => void | Promise<void>196
}198
/** Dependency overrides for an isolated generation runtime. */199
export type JsonlGenerationRuntimeOverrides = Partial<Omit<JsonlGenerationInternals, 'fs'>> & {200
readonly fs?: Partial<GenerationFileSystem>201
}203
/** Bound generation operations used by production defaults and deterministic tests. */204
export interface JsonlGenerationRuntime {205
readStable(path: string, signal?: AbortSignal): Promise<StablePhysicalFile>206
prepare(options: PrepareJsonlMigrationOptions): Promise<PreparedJsonlMigration>207
verify(208
path: string,209
compression: JsonlCompression,210
expectedId: string,211
expectedEventCount: number,212
expectedPrefix?: JsonlExpectedPrefix,213
): Promise<JsonlVerifiedGeneration>214
}216
const defaultFileSystem: GenerationFileSystem = {217
open: (path, flags, mode) => fsOpen(path, flags, mode),218
readFile: (path, signal) => fsReadFile(path, signal === undefined ? undefined : { signal }),219
readdir: path => fsReaddir(path),220
stat: path => fsStat(path, { bigint: true }),221
lstat: path => fsLstat(path),222
link: fsLink,223
rm: path => fsRm(path, { force: true }),224
}226
const defaultInternals: JsonlGenerationInternals = {227
fs: defaultFileSystem,228
randomToken: () => randomBytes(8).toString('hex'),229
platform: process.platform,230
publishNewWin32: publishNewFileWin32,231
barrier: () => {},232
}234
function isEEXIST(error: unknown): boolean {235
return (error as NodeJS.ErrnoException | null)?.code === 'EEXIST'236
}238
/** Whether a filesystem-owned failure should retain its original errno and path. */239
function isErrnoException(error: unknown): error is NodeJS.ErrnoException {240
return typeof (error as NodeJS.ErrnoException | null)?.code === 'string'241
}243
function identity(value: JsonlPhysicalIdentity): string {244
return [value.dev, value.ino, value.size, value.mtimeNs, value.ctimeNs].join(':')245
}247
/**248
* Read one stable revision of a JSONL file with a single retry. If an append249
* overlaps both reads, return the second read's committed pre-read prefix250
* instead of starving behind a continuous writer.251
* @param path - the generation file to read.252
* @param signal - optional cancellation for the stat/read work.253
* @returns the stable bytes (or the committed prefix) and their stat identity.254
*/255
export async function readStableJsonlFile(256
path: string,257
signal?: AbortSignal,258
): Promise<StablePhysicalFile> {259
return defaultGenerationRuntime.readStable(path, signal)260
}262
async function readStableSnapshot(263
path: string,264
signal: AbortSignal | undefined,265
fs: GenerationFileSystem,266
): Promise<StablePhysicalFile> {267
signal?.throwIfAborted()268
let before = await fs.stat(path)269
for (let attempt = 0; ; attempt += 1) {270
const bytes = await fs.readFile(path, signal)271
signal?.throwIfAborted()272
const after = await fs.stat(path)273
if (identity(before) === identity(after)) {274
signal?.throwIfAborted()275
return { bytes, identity: after }276
}277
if (attempt === 1) {278
return { bytes: bytes.subarray(0, Number(before.size)), identity: before }279
}280
before = after281
}282
}284
/** Parse the version discriminator without validating any version-specific field. */285
function storedVersion(header: unknown): number {286
if (typeof header !== 'object' || header === null || Array.isArray(header)) {287
throw new Error('corrupt session log: first line is not a JSON object')288
}289
const version = (header as { version?: unknown }).version290
if (!Number.isSafeInteger(version) || (version as number) < 0 || Object.is(version, -0)) {291
throw new Error('corrupt session log: header version is not a non-negative safe integer')292
}293
return version as number294
}296
function parseJson(text: string, subject: string): unknown {297
try {298
return JSON.parse(text)299
} catch (error) {300
throw new Error(`corrupt session log: ${subject} is not valid JSON`, { cause: error })301
}302
}304
/** Incremental JSONL parser that retains only one cross-frame record fragment. */305
class MigratingJsonlRows {306
private fragments: Buffer[] = []307
private fragmentBytes = 0308
private rowIndex = 0309
private issue: Error | undefined311
constructor(private readonly restore: SessionFormatRestore) {}313
/** Consume plaintext bytes following the independently decoded header. */314
write(chunk: Buffer): void {315
/* jscpd:ignore-start -- migration parsing and readable-log scanning own different recovery and byte-accounting state. */316
let lineStart = 0317
for (318
let newline = chunk.indexOf(0x0A);319
newline !== -1;320
newline = chunk.indexOf(0x0A, lineStart)321
) {322
const fragment = chunk.subarray(lineStart, newline)323
let line = fragment324
if (this.fragments.length > 0) {325
if (fragment.length > 0) this.fragments.push(fragment)326
line = Buffer.concat(this.fragments, this.fragmentBytes + fragment.length)327
this.fragments = []328
this.fragmentBytes = 0329
}330
this.consume(line)331
lineStart = newline + 1332
}333
if (lineStart < chunk.length) {334
const fragment = Buffer.from(chunk.subarray(lineStart))335
this.fragments.push(fragment)336
this.fragmentBytes += fragment.length337
}338
/* jscpd:ignore-end */339
}341
/** Refuse a record fragment left by structurally complete Zstandard frames. */342
assertCompleteFramesEndOnRecord(): void {343
if (this.fragments.length > 0) {344
throw new Error('corrupt Zstandard session log: complete frame contains a torn JSONL record')345
}346
}348
finish(): SessionFormatArtifact {349
return this.restore.finish()350
}352
private consume(line: Buffer): void {353
const index = this.rowIndex354
this.rowIndex += 1355
let row: unknown356
try {357
row = parseJson(line.toString('utf8'), `row ${index + 1}`)358
} catch (error: unknown) {359
this.issue ??= asError(error)360
return361
}362
if (this.issue !== undefined) {363
if (typeof row === 'object' && row !== null364
&& (row as { type?: unknown }).type === 'turn/end') throw this.issue365
return366
}367
this.restore.decodeRow(row)368
}369
}371
interface StartedMigrationStream {372
readonly parser: MigratingJsonlRows373
}375
async function startMigrationStream(376
headerRecord: Buffer,377
sourceVersion: number,378
format: Pick<JsonlGenerationFormatAdapter, 'createRestore'>,379
validateHistoricalHeader?: PrepareJsonlMigrationOptions['validateHistoricalHeader'],380
): Promise<StartedMigrationStream> {381
const value = parseJson(headerRecord.subarray(0, -1).toString('utf8'), 'header line')382
const version = storedVersion(value)383
if (version !== sourceVersion) {384
throw new Error(`resolved JSONL source filename identifies v${sourceVersion}, but its header identifies v${version}`)385
}386
const header = value as Record<string, unknown>387
const validation = validateHistoricalHeader?.(header)388
if (validation !== undefined) await validation389
const stream = format.createRestore(header)390
return { parser: new MigratingJsonlRows(stream) }391
}393
async function consumeMigrationBytes(394
rows: MigratingJsonlRows,395
chunks: Iterable<Buffer>,396
signal?: AbortSignal,397
): Promise<void> {398
signal?.throwIfAborted()399
let yieldDeadline = performance.now() + MIGRATION_DECODE_YIELD_INTERVAL_MS400
for (const bytes of chunks) {401
for (let offset = 0; offset < bytes.length; offset += MIGRATION_WORK_CHUNK_BYTES) {402
rows.write(bytes.subarray(offset, offset + MIGRATION_WORK_CHUNK_BYTES))403
if (performance.now() < yieldDeadline) continue404
await scheduler.yield()405
signal?.throwIfAborted()406
yieldDeadline = performance.now() + MIGRATION_DECODE_YIELD_INTERVAL_MS407
}408
}409
}411
async function decodeStreamingMigration(412
bytes: Buffer,413
compression: JsonlCompression,414
sourceVersion: number,415
format: Pick<JsonlGenerationFormatAdapter, 'createRestore'>,416
validateHistoricalHeader: PrepareJsonlMigrationOptions['validateHistoricalHeader'],417
signal?: AbortSignal,418
): Promise<SessionFormatArtifact> {419
signal?.throwIfAborted()420
if (compression === 'none') {421
const headerEnd = bytes.indexOf(0x0A)422
if (headerEnd === -1) throw new Error('empty or header-less session log')423
const stream = await startMigrationStream(424
bytes.subarray(0, headerEnd + 1),425
sourceVersion,426
format,427
validateHistoricalHeader,428
)429
signal?.throwIfAborted()430
const bodyEnd = bytes.lastIndexOf(0x0A)431
if (bodyEnd > headerEnd) {432
await consumeMigrationBytes(433
stream.parser,434
[bytes.subarray(headerEnd + 1, bodyEnd + 1)],435
signal,436
)437
}438
return stream.parser.finish()439
}441
const { frames, tornStart } = scanZstdFrames(bytes)442
if (frames.length === 0) throw new Error('empty or header-less Zstandard session log')443
const decoder = createZstdFrameDecoder()444
try {445
const decoded = decoder.decode(bytes, frames)446
const first = decoded.next()447
/* v8 ignore next -- a non-empty structural frame list yields once or throws. */448
if (first.done) throw new Error('empty or header-less Zstandard session log')449
assertIndependentHeaderFrame(first.value)450
const stream = await startMigrationStream(451
first.value,452
sourceVersion,453
format,454
validateHistoricalHeader,455
)456
signal?.throwIfAborted()457
await consumeMigrationBytes(stream.parser, decoded, signal)458
stream.parser.assertCompleteFramesEndOnRecord()459
if (tornStart !== undefined) {460
let recovered: Buffer = Buffer.alloc(0)461
try {462
recovered = await decompressZstdPrefix(bytes.subarray(tornStart))463
} catch {464
/* v8 ignore next -- decoder failure plus concurrent abort is timing-dependent. */465
if (signal?.aborted) signal.throwIfAborted()466
}467
signal?.throwIfAborted()468
const newline = recovered.lastIndexOf(0x0A)469
if (newline !== -1) {470
await consumeMigrationBytes(471
stream.parser,472
[recovered.subarray(0, newline + 1)],473
signal,474
)475
}476
}477
return stream.parser.finish()478
} finally {479
decoder.close()480
}481
}483
/**484
* Read and validate one complete current generation for an isolated verifier.485
* @param path - staged or competing current-generation path.486
* @param compression - configured physical encoding.487
* @param expectedId - Session identity expected in the header.488
* @param expectedEventCount - exact logical event count expected after decoding.489
* @param expectedPrefix - verified migration prefix; an append tail may be present and is not validated.490
* @returns stable physical identity and digest for publication comparison.491
*/492
export async function verifyJsonlCurrentGeneration(493
path: string,494
compression: JsonlCompression,495
expectedId: string,496
expectedEventCount: number,497
expectedPrefix?: JsonlExpectedPrefix,498
): Promise<JsonlVerifiedGeneration> {499
return defaultGenerationRuntime.verify(path, compression, expectedId, expectedEventCount, expectedPrefix)500
}502
async function verifyCurrentGeneration(503
path: string,504
compression: JsonlCompression,505
expectedId: string,506
expectedEventCount: number,507
fs: GenerationFileSystem,508
expectedPrefix?: JsonlExpectedPrefix,509
): Promise<JsonlVerifiedGeneration> {510
const before = await fs.stat(path)511
const bytes = await fs.readFile(path)512
const after = await fs.stat(path)513
if (expectedPrefix !== undefined) {514
if (bytes.length < expectedPrefix.bytes) {515
throw new Error('target bytes are shorter than the migrated generation')516
}517
const digest = createHash('sha256').update(bytes.subarray(0, expectedPrefix.bytes)).digest('hex')518
if (digest !== expectedPrefix.digest) {519
throw new Error('target bytes do not begin with the migrated generation')520
}521
return { identity: after, bytes: expectedPrefix.bytes, digest }522
}523
if (identity(before) !== identity(after)) {524
throw new Error('current session generation changed during verification')525
}526
const snapshot = { bytes, identity: after }527
const generation = decodeCurrentGeneration(snapshot.bytes, compression)528
validateStoredEvents(generation.meta, generation.events, { kind: 'jsonl', path })529
if (generation.meta.id !== expectedId) {530
throw new Error(`current session generation contains id "${generation.meta.id}", expected "${expectedId}"`)531
}532
if (generation.events.length !== expectedEventCount) {533
throw new Error(534
`current session generation contains ${generation.events.length} events, expected ${expectedEventCount}`,535
)536
}537
Session.fromRestore(538
generation.meta.id,539
generation.events,540
generation.meta,541
generation.inheritedEventCount,542
'detached',543
currentSessionMessageProjections,544
)545
assertCurrentAssistantStreams(generation.events)546
return {547
identity: snapshot.identity,548
bytes: snapshot.bytes.length,549
digest: createHash('sha256').update(snapshot.bytes).digest('hex'),550
}551
}553
/** Fully replay embedded streams only inside isolated current-generation verification. */554
function assertCurrentAssistantStreams(events: readonly SessionEvent[]): void {555
for (const [index, event] of events.entries()) {556
if (event.type !== 'assistant/message' && event.type !== 'assistant/attempt') continue557
const assembler = new BlockAssembler()558
let timed: ReturnType<typeof expandAssistantStream>559
try {560
timed = expandAssistantStream(event.data.stream)561
for (const member of timed) assembler.push(member.chunk)562
} catch (error: unknown) {563
throw new Error(`seed ${event.type} at index ${index} has an invalid embedded stream`, { cause: error })564
}565
if (event.type === 'assistant/attempt' || timed.length === 0) continue566
const content = event.data.interrupted === true ? assembler.interruptedBlocks() : assembler.blocks()567
if (!isDeepStrictEqual(event.data.message.content, content)) {568
throw new Error(`seed assistant/message at index ${index} content disagrees with its embedded stream`)569
}570
if (!isDeepStrictEqual(event.data.usage, assembler.usage)) {571
throw new Error(`seed assistant/message at index ${index} usage disagrees with its embedded stream`)572
}573
if (!isDeepStrictEqual(event.data.message.source.replayState, assembler.replayState)) {574
throw new Error(`seed assistant/message at index ${index} replay state disagrees with its embedded stream`)575
}576
}577
}579
function decodeCurrentGeneration(580
bytes: Buffer,581
compression: JsonlCompression,582
): ReturnType<SessionLogScanner['finish']> {583
if (compression === 'none') {584
const headerEnd = bytes.indexOf(0x0A)585
if (headerEnd === -1) throw new Error('empty or header-less session log')586
const scanner = new SessionLogScanner(bytes.subarray(0, headerEnd + 1), 'strict')587
scanner.write(bytes.subarray(headerEnd + 1))588
return finishCurrentGenerationScan(scanner)589
}590
const { frames, tornStart } = scanZstdFrames(bytes)591
if (frames.length === 0) throw new Error('empty or header-less Zstandard session log')592
if (tornStart !== undefined) throw new Error('current session generation has a torn physical tail')593
const decoder = createZstdFrameDecoder()594
try {595
const plaintext = decoder.decode(bytes, frames)596
const header = plaintext.next()597
/* v8 ignore next -- a non-empty structural frame list yields once or throws. */598
if (header.done) throw new Error('empty or header-less Zstandard session log')599
assertIndependentHeaderFrame(header.value)600
const scanner = new SessionLogScanner(header.value, 'strict')601
for (const chunk of plaintext) scanner.write(chunk)602
return finishCurrentGenerationScan(scanner)603
} finally {604
decoder.close()605
}606
}608
function finishCurrentGenerationScan(609
scanner: SessionLogScanner,610
): ReturnType<SessionLogScanner['finish']> {611
const inputBytes = scanner.checkpoint().inputBytes612
const decoded = scanner.finish()613
if (decoded.committedBytes !== inputBytes) throw new Error('current session generation has a torn physical tail')614
return decoded615
}617
function stringifyJson(value: unknown, subject: string): string {618
let text: unknown619
try {620
text = JSON.stringify(value)621
} catch (error) {622
throw new Error(`${subject} is not lossless JSON`, { cause: error })623
}624
if (typeof text !== 'string') throw new Error(`${subject} is not lossless JSON`)625
return text626
}628
function assertIndependentHeaderFrame(plaintext: Buffer): void {629
if (plaintext.length === 0 || plaintext.indexOf(0x0A) !== plaintext.length - 1) {630
throw new Error('corrupt Zstandard session log: first frame is not exactly one header line')631
}632
}634
function assertGenerationPaths(635
sourcePath: string,636
sourceVersion: number,637
currentPath: string,638
currentVersion: number,639
compression: JsonlCompression,640
): string {641
const expectedSource = generationLogFilename(sourceVersion, compression)642
const expectedCurrent = generationLogFilename(currentVersion, compression)643
if (basename(sourcePath) !== expectedSource) {644
throw new Error(`resolved JSONL source path must end with "${expectedSource}": ${sourcePath}`)645
}646
if (basename(currentPath) !== expectedCurrent) {647
throw new Error(`current JSONL generation path must end with "${expectedCurrent}": ${currentPath}`)648
}649
if (dirname(sourcePath) !== dirname(currentPath)) {650
throw new Error('source and current JSONL generations must share one Session directory')651
}652
return logSuffix(compression)653
}655
async function syncDirectory(path: string, internals: JsonlGenerationInternals): Promise<void> {656
/* v8 ignore next -- Windows namespace operations request write-through directly. */657
if (internals.platform === 'win32') return658
const handle = await internals.fs.open(path, 'r')659
try {660
await handle.sync()661
} finally {662
await handle.close()663
}664
}666
interface StreamedMigrationStage {667
readonly path: string668
readonly bytes: number669
readonly digest: string670
}672
/** Produce bounded JSONL chunks while yielding between main-thread encoding slices. */673
async function* encodeMigrationRows(674
artifact: SessionFormatArtifact,675
format: JsonlGenerationFormatAdapter,676
signal?: AbortSignal,677
): AsyncGenerator<Buffer, void, void> {678
signal?.throwIfAborted()679
let lines: string[] = []680
let bytes = 0681
for (const value of artifact.events) {682
const line = `${stringifyJson(format.encodeEvent(value), `migrated Session event ${value.seq}`)}\n`683
const lineBytes = Buffer.byteLength(line)684
if (bytes > 0 && bytes + lineBytes > MIGRATION_WORK_CHUNK_BYTES) {685
yield Buffer.from(lines.join(''))686
await scheduler.yield()687
signal?.throwIfAborted()688
lines = []689
bytes = 0690
}691
lines.push(line)692
bytes += lineBytes693
}694
yield Buffer.from(lines.join(''))695
}697
async function writeMigrationChunks(698
chunks: AsyncIterable<Buffer>,699
write: (chunk: Buffer) => Promise<void>,700
): Promise<void> {701
let pending: Buffer[] = []702
let bytes = 0703
for await (const chunk of chunks) {704
pending.push(chunk)705
bytes += chunk.length706
if (bytes < MIGRATION_WRITE_CHUNK_BYTES) continue707
await write(pending.length === 1 ? pending[0] as Buffer : Buffer.concat(pending, bytes))708
pending = []709
bytes = 0710
}711
if (bytes > 0) await write(pending.length === 1 ? pending[0] as Buffer : Buffer.concat(pending, bytes))712
}714
/** Encode directly into one synced stage without a whole-artifact row or byte buffer. */715
async function writeSyncedTemp(716
currentPath: string,717
suffix: string,718
compression: JsonlCompression,719
artifact: SessionFormatArtifact,720
format: JsonlGenerationFormatAdapter,721
signal: AbortSignal | undefined,722
internals: JsonlGenerationInternals,723
): Promise<StreamedMigrationStage> {724
signal?.throwIfAborted()725
let path: string726
let handle: FileHandle727
for (;;) {728
path = join(dirname(currentPath), `session.migration.${internals.randomToken()}${suffix}.tmp`)729
try {730
handle = await internals.fs.open(path, 'wx', 0o600)731
break732
} catch (error) {733
if (isEEXIST(error)) continue734
throw error735
}736
}737
const hash = createHash('sha256')738
let bytes = 0739
const write = async (chunk: Buffer): Promise<void> => {740
await handle.writeFile(chunk)741
hash.update(chunk)742
bytes += chunk.length743
}744
let failure: unknown745
try {746
const headerValue = format.encodeHeader(artifact.header, artifact.inheritedEventCount)747
const header = Buffer.from(`${stringifyJson(headerValue, 'migrated session header')}\n`)748
await write(compression === 'zstd' ? await compressZstdFrame(header) : header)749
if (artifact.events.length > 0) {750
const rows = encodeMigrationRows(artifact, format, signal)751
if (compression === 'none') {752
await writeMigrationChunks(rows, write)753
} else {754
await new Promise<void>((resolve, reject) => {755
pipeline(756
Readable.from(rows, { objectMode: false, highWaterMark: MIGRATION_WORK_CHUNK_BYTES }),757
createZstdCompress(ZSTD_CHECKSUM_OPTIONS),758
async (source) => { await writeMigrationChunks(source as AsyncIterable<Buffer>, write) },759
(error: Error | null | undefined) => {760
if (error instanceof Error) reject(error)761
else resolve()762
},763
)764
})765
}766
}767
signal?.throwIfAborted()768
await handle.sync()769
} catch (error: unknown) {770
failure = error771
}772
try {773
await handle.close()774
} catch (error: unknown) {775
failure = failure === undefined776
? error777
: new AggregateError([failure, error], `failed to write and close migration stage "${path}"`)778
}779
if (failure !== undefined) {780
const writeError = failure instanceof Error781
? failure782
: new Error('migration stage write failed with a non-Error rejection', { cause: failure })783
await removeTemporary(path, writeError, internals)784
throw writeError785
}786
return { path, bytes, digest: hash.digest('hex') }787
}789
/** Remove one temporary file without hiding the operation failure that made it disposable. */790
async function removeTemporary(791
path: string,792
primaryFailure: unknown,793
internals: JsonlGenerationInternals,794
): Promise<void> {795
try {796
await internals.fs.rm(path)797
} catch (cleanupFailure: unknown) {798
throw new AggregateError(799
[primaryFailure, cleanupFailure],800
`failed to clean migration temporary "${path}" after an earlier failure`,801
)802
}803
}805
/** Remove a redundant stage after the target has been validated as committed. */806
async function removeCommittedTemporary(807
path: string,808
internals: JsonlGenerationInternals,809
): Promise<void> {810
try {811
await internals.fs.rm(path)812
} catch {813
// The validated target owns the committed bytes; a redundant link cannot turn success into failure.814
}815
}817
async function publishCurrentExclusive(818
staged: string,819
currentPath: string,820
internals: JsonlGenerationInternals,821
): Promise<boolean> {822
if (internals.platform === 'win32') {823
try {824
await internals.publishNewWin32(staged, currentPath)825
return true826
} catch (error) {827
/* v8 ignore else -- native helper tests own non-collision Win32 failures. */828
if (isEEXIST(error)) return false829
/* v8 ignore next -- the filesystem error is already complete. */830
throw error831
}832
}833
try {834
await internals.fs.link(staged, currentPath)835
} catch (error) {836
/* v8 ignore else -- a non-collision filesystem error propagates unchanged. */837
if (isEEXIST(error)) return false838
/* v8 ignore next -- the filesystem error is already complete. */839
throw error840
}841
await syncDirectory(dirname(currentPath), internals)842
return true843
}845
function asError(error: unknown): Error {846
return error instanceof Error ? error : new Error('current-generation validation failed with a non-Error rejection', {847
cause: error,848
})849
}851
async function inspectExpectedCurrent<T>(852
currentPath: string,853
internals: JsonlGenerationInternals,854
inspect: () => Promise<T>,855
): Promise<T> {856
try {857
const expectedName = basename(currentPath)858
const names = await internals.fs.readdir(dirname(currentPath))859
if (!names.includes(expectedName)) {860
const noncanonical = names.find(name => name.toLowerCase() === expectedName.toLowerCase())861
if (noncanonical !== undefined) {862
throw new Error(`target resolves to noncanonical directory entry "${noncanonical}"`)863
}864
}865
const info = await internals.fs.lstat(currentPath)866
if (info.isSymbolicLink() || !info.isFile()) {867
throw new Error(`target is a ${info.isSymbolicLink() ? 'symbolic link' : 'non-regular file'}`)868
}869
return await inspect()870
} catch (error: unknown) {871
if (isErrnoException(error)) throw error872
throw new JsonlGenerationTargetConflictError(currentPath, asError(error))873
}874
}876
function withOverrides(overrides: JsonlGenerationRuntimeOverrides): JsonlGenerationInternals {877
return {878
...defaultInternals,879
...overrides,880
fs: { ...defaultFileSystem, ...overrides.fs },881
}882
}884
async function publishPreparedMigration(885
options: PrepareJsonlMigrationOptions,886
suffix: string,887
artifact: SessionFormatArtifact,888
sourceIdentity: JsonlPhysicalIdentity,889
internals: JsonlGenerationInternals,890
): Promise<JsonlPhysicalIdentity> {891
await scheduler.yield()892
const { sourcePath, currentPath, compression, verifyCurrentFile } = options893
const eventCount = artifact.events.length894
let staged = await writeSyncedTemp(currentPath, suffix, compression, artifact, options.format, undefined, internals)895
try {896
const verifiedStage = await verifyCurrentFile(897
staged.path,898
compression,899
artifact.header.id,900
eventCount,901
)902
if (verifiedStage.bytes !== staged.bytes || verifiedStage.digest !== staged.digest) {903
throw new Error('staged session generation changed during verification')904
}905
await internals.barrier('before-source-check', 1)906
await options.validateRelatedSources?.()907
const beforePublish = await internals.fs.stat(sourcePath)908
if (identity(beforePublish) !== identity(sourceIdentity)) {909
throw new JsonlGenerationSourceChangedError(sourcePath)910
}911
const published = await publishCurrentExclusive(staged.path, currentPath, internals)912
if (published && internals.platform === 'win32') staged = { ...staged, path: '' }913
await internals.barrier('after-publication', 1)914
let currentIdentity: JsonlPhysicalIdentity915
if (published) {916
if (staged.path !== '') {917
await removeCommittedTemporary(staged.path, internals)918
staged = { ...staged, path: '' }919
}920
currentIdentity = await internals.fs.stat(currentPath)921
} else {922
const winner = await inspectExpectedCurrent(currentPath, internals, async () => {923
const candidate = await verifyCurrentFile(924
currentPath,925
compression,926
artifact.header.id,927
eventCount,928
staged,929
)930
if (candidate.bytes !== staged.bytes || candidate.digest !== staged.digest) {931
throw new Error('target bytes differ from the migrated generation')932
}933
return candidate934
})935
currentIdentity = winner.identity936
await removeCommittedTemporary(staged.path, internals)937
staged = { ...staged, path: '' }938
}939
return currentIdentity940
} catch (error: unknown) {941
if (staged.path !== '') await removeTemporary(staged.path, error, internals)942
throw error943
}944
}946
async function prepareMigration(947
options: PrepareJsonlMigrationOptions,948
internals: JsonlGenerationInternals,949
): Promise<PreparedJsonlMigration> {950
const { sourcePath, sourceVersion, currentPath, compression, format, signal } = options951
const suffix = assertGenerationPaths(952
sourcePath,953
sourceVersion,954
currentPath,955
format.currentVersion,956
compression,957
)958
if (sourceVersion >= format.currentVersion) {959
throw new Error(`migration preparation requires a historical source, got v${sourceVersion}`)960
}961
const source = await readStableSnapshot(sourcePath, signal, internals.fs)962
let artifact: SessionFormatArtifact963
try {964
artifact = await decodeStreamingMigration(965
source.bytes,966
compression,967
sourceVersion,968
format,969
options.validateHistoricalHeader,970
signal,971
)972
} catch (error: unknown) {973
if (format.isUnsupportedMigrationError?.(error) === true) {974
throw new JsonlGenerationUnsupportedMigrationError(sourceVersion, error)975
}976
throw error977
}978
if (artifact.header.version !== format.currentVersion) {979
throw new Error(`format migration returned v${artifact.header.version}, expected v${format.currentVersion}`)980
}981
await options.validateRelatedSources?.()982
const sourceIdentity = source.identity983
let publication: Promise<JsonlPhysicalIdentity> | undefined984
return {985
sourceIdentity,986
artifact,987
publish() {988
if (publication === undefined) {989
publication = publishPreparedMigration(990
options,991
suffix,992
artifact,993
sourceIdentity,994
internals,995
)996
}997
return publication998
},999
}1000
}1002
/**1003
* Decode and migrate one historical generation without writing its successor.1004
* @param options - resolved source, current target, format adapter, and load cancellation.1005
* @returns the current artifact and an idempotent explicit publication operation.1006
*/1007
export function prepareJsonlMigration(1008
options: PrepareJsonlMigrationOptions,1009
): Promise<PreparedJsonlMigration> {1010
return defaultGenerationRuntime.prepare(options)1011
}1013
/**1014
* Create one generation runtime with fixed filesystem and publication dependencies.1015
* @param overrides - deterministic filesystem, platform, and race dependencies.1016
* @returns bound generation operations.1017
*/1018
export function createJsonlGenerationRuntime(1019
overrides: JsonlGenerationRuntimeOverrides = {},1020
): JsonlGenerationRuntime {1021
const internals = withOverrides(overrides)1022
return {1023
readStable: (path, signal) => readStableSnapshot(path, signal, internals.fs),1024
prepare: options => prepareMigration(options, internals),1025
verify: (path, compression, expectedId, expectedEventCount, expectedPrefix) => verifyCurrentGeneration(1026
path, compression, expectedId, expectedEventCount, internals.fs, expectedPrefix,1027
),1028
}1029
}1031
const defaultGenerationRuntime = createJsonlGenerationRuntime()1033
/**1034
* Read one stable source through the shared streaming parser without publishing a generation.1035
* @param path - selected source generation path.1036
* @param version - physical source version identified by its filename.1037
* @param compression - source encoding.1038
* @param format - codec/restore factory, independent of current-generation publication.1039
* @param signal - cancellation observed during source reads and decode yields.1040
* @returns decoded artifact and physical source identity for later revalidation.1041
* @throws SessionFormatError for physical decoding failures; storage, cancellation, and unsupported migration errors retain their category.1042
*/1043
export async function readDecodedJsonlSource(1044
path: string,1045
version: number,1046
compression: JsonlCompression,1047
format: Pick<JsonlGenerationFormatAdapter, 'createRestore'>,1048
signal?: AbortSignal,1049
): Promise<{ artifact: SessionFormatArtifact; identity: JsonlPhysicalIdentity }> {1050
const source = await readStableJsonlFile(path, signal)1051
let artifact: SessionFormatArtifact1052
try {1053
artifact = await decodeStreamingMigration(source.bytes, compression, version, format, undefined, signal)1054
} catch (error: unknown) {1055
if (signal?.aborted || error instanceof SessionFormatError) throw error1056
throw new SessionFormatError(String(error), { cause: error })1057
}1058
return { artifact, identity: source.identity }1059
}