返回源码地图

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

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

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

1/**
2 * Durable whole-generation publication for JSONL Session artifacts.
3 *
4 * Format packages transform parsed JSON values. This module owns the physical
5 * encoding, exact source identity, immutable generation files, and exclusive
6 * current-generation publication for both configured JSONL suffixes.
7 * @module @deepseek-ai/dsh-session-persistence-jsonl/generation
8 */
9
10import { currentSessionMessageProjections } from '@deepseek-ai/dsh-session-format-catalog/message-projections'
11import { createHash, randomBytes } from 'node:crypto'
12import {
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'
22import { basename, dirname, join } from 'node:path'
23import { performance } from 'node:perf_hooks'
24import { pipeline, Readable } from 'node:stream'
25import { scheduler } from 'node:timers/promises'
26import { isDeepStrictEqual } from 'node:util'
27import { constants, createZstdCompress } from 'node:zlib'
28import { Session } from '@deepseek-ai/dsh-session'
29import type { SessionEvent } from '@deepseek-ai/dsh-session'
30import { BlockAssembler, expandAssistantStream } from '@deepseek-ai/dsh-llm'
31import { SessionFormatError } from '@deepseek-ai/dsh-session-format'
32import type {
33 SessionFormatArtifact,
34 SessionFormatJsonValue,
35 SessionFormatRestore,
36} from '@deepseek-ai/dsh-session-format'
37import { validateStoredEvents } from '@deepseek-ai/dsh-session-persistence'
38import type { JsonlCompression } from './format.ts'
39import { generationLogFilename, logSuffix, SessionLogScanner } from './format.ts'
40import { publishNewFileWin32 } from './win32.ts'
41import {
42 compressZstdFrame,
43 createZstdFrameDecoder,
44 decompressZstdPrefix,
45 scanZstdFrames,
46} from './zstd.ts'
47
48/** Internal scheduling bounds: preserve old decode cadence and cap each synchronous encode slice. */
49const MIGRATION_DECODE_YIELD_INTERVAL_MS = 500
50const MIGRATION_WORK_CHUNK_BYTES = 1024 * 1024
51const MIGRATION_WRITE_CHUNK_BYTES = 4 * 1024 * 1024
52const ZSTD_CHECKSUM_OPTIONS = {
53 chunkSize: MIGRATION_WORK_CHUNK_BYTES,
54 params: { [constants.ZSTD_c_checksumFlag]: 1 },
55}
56
57/** Pure adapter between backend-owned JSONL framing and the format catalog. */
58export interface JsonlGenerationFormatAdapter {
59 readonly currentVersion: number
60 /** Create the single-pass codec and migration state for a historical header. */
61 createRestore(header: Record<string, unknown>): SessionFormatRestore
62 /** Encode one current header record without materializing body rows. */
63 encodeHeader(header: SessionFormatArtifact['header'], inheritedEventCount: number): SessionFormatJsonValue
64 /** Encode one current event record. */
65 encodeEvent(event: SessionFormatArtifact['events'][number]): SessionFormatJsonValue
66 /** Classify a supported-version artifact that policy refuses to migrate. */
67 isUnsupportedMigrationError?(error: unknown): error is Error
68}
69
70/** Inputs for preparing one historical generation and publishing its current successor later. */
71export 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: string
76 /** Version selected from the source filename and independently checked against its header. */
77 readonly sourceVersion: number
78 /** Canonical filename for `format.currentVersion` in the same Session directory. */
79 readonly currentPath: string
80 readonly compression: JsonlCompression
81 readonly format: JsonlGenerationFormatAdapter
82 /** 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?: AbortSignal
96}
97
98/** Small physical identity returned by an isolated generation verifier. */
99export interface JsonlVerifiedGeneration {
100 readonly identity: JsonlPhysicalIdentity
101 readonly bytes: number
102 readonly digest: string
103}
104
105/** Physical byte prefix already proven to be a valid complete generation. */
106export interface JsonlExpectedPrefix {
107 readonly bytes: number
108 readonly digest: string
109}
110
111/** A historical source changed after its single decode and migration pass. */
112export class JsonlGenerationSourceChangedError extends Error {
113 override readonly name = 'JsonlGenerationSourceChangedError'
114
115 /** @param path - historical generation whose revision changed. */
116 constructor(readonly path: string) {
117 super(`historical session generation changed during migration: "${path}"`)
118 }
119}
120
121/** Current logical state prepared independently from durable publication. */
122export interface PreparedJsonlMigration {
123 readonly sourceIdentity: JsonlPhysicalIdentity
124 readonly artifact: SessionFormatArtifact
125 /** Encode, verify, and exclusively publish once; every call shares the same success or failure. */
126 publish(): Promise<JsonlPhysicalIdentity>
127}
128
129/** A historical artifact is intact, but the format edge refuses its contents. */
130export class JsonlGenerationUnsupportedMigrationError extends Error {
131 override readonly name = 'JsonlGenerationUnsupportedMigrationError'
132
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}
144
145/** A current-generation filename already names different or invalid bytes. */
146export class JsonlGenerationTargetConflictError extends Error {
147 override readonly name = 'JsonlGenerationTargetConflictError'
148
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}
160
161/** Stat identity captured together with exact generation bytes. */
162export interface JsonlPhysicalIdentity {
163 readonly dev: bigint
164 readonly ino: bigint
165 readonly size: bigint
166 readonly mtimeNs: bigint
167 readonly ctimeNs: bigint
168}
169
170/** Exact bytes of one stable file revision together with the stat identity that proved it stable. */
171export interface StablePhysicalFile {
172 readonly bytes: Buffer
173 readonly identity: JsonlPhysicalIdentity
174}
175
176interface 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}
185
186type GenerationBarrierPhase =
187 | 'before-source-check'
188 | 'after-publication'
189
190interface JsonlGenerationInternals {
191 readonly fs: GenerationFileSystem
192 readonly randomToken: () => string
193 readonly platform: NodeJS.Platform
194 readonly publishNewWin32: typeof publishNewFileWin32
195 readonly barrier: (phase: GenerationBarrierPhase, attempt: number) => void | Promise<void>
196}
197
198/** Dependency overrides for an isolated generation runtime. */
199export type JsonlGenerationRuntimeOverrides = Partial<Omit<JsonlGenerationInternals, 'fs'>> & {
200 readonly fs?: Partial<GenerationFileSystem>
201}
202
203/** Bound generation operations used by production defaults and deterministic tests. */
204export 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}
215
216const 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}
225
226const defaultInternals: JsonlGenerationInternals = {
227 fs: defaultFileSystem,
228 randomToken: () => randomBytes(8).toString('hex'),
229 platform: process.platform,
230 publishNewWin32: publishNewFileWin32,
231 barrier: () => {},
232}
233
234function isEEXIST(error: unknown): boolean {
235 return (error as NodeJS.ErrnoException | null)?.code === 'EEXIST'
236}
237
238/** Whether a filesystem-owned failure should retain its original errno and path. */
239function isErrnoException(error: unknown): error is NodeJS.ErrnoException {
240 return typeof (error as NodeJS.ErrnoException | null)?.code === 'string'
241}
242
243function identity(value: JsonlPhysicalIdentity): string {
244 return [value.dev, value.ino, value.size, value.mtimeNs, value.ctimeNs].join(':')
245}
246
247/**
248 * Read one stable revision of a JSONL file with a single retry. If an append
249 * overlaps both reads, return the second read's committed pre-read prefix
250 * 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 */
255export async function readStableJsonlFile(
256 path: string,
257 signal?: AbortSignal,
258): Promise<StablePhysicalFile> {
259 return defaultGenerationRuntime.readStable(path, signal)
260}
261
262async 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 = after
281 }
282}
283
284/** Parse the version discriminator without validating any version-specific field. */
285function 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 }).version
290 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 number
294}
295
296function 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}
303
304/** Incremental JSONL parser that retains only one cross-frame record fragment. */
305class MigratingJsonlRows {
306 private fragments: Buffer[] = []
307 private fragmentBytes = 0
308 private rowIndex = 0
309 private issue: Error | undefined
310
311 constructor(private readonly restore: SessionFormatRestore) {}
312
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 = 0
317 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 = fragment
324 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 = 0
329 }
330 this.consume(line)
331 lineStart = newline + 1
332 }
333 if (lineStart < chunk.length) {
334 const fragment = Buffer.from(chunk.subarray(lineStart))
335 this.fragments.push(fragment)
336 this.fragmentBytes += fragment.length
337 }
338 /* jscpd:ignore-end */
339 }
340
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 }
347
348 finish(): SessionFormatArtifact {
349 return this.restore.finish()
350 }
351
352 private consume(line: Buffer): void {
353 const index = this.rowIndex
354 this.rowIndex += 1
355 let row: unknown
356 try {
357 row = parseJson(line.toString('utf8'), `row ${index + 1}`)
358 } catch (error: unknown) {
359 this.issue ??= asError(error)
360 return
361 }
362 if (this.issue !== undefined) {
363 if (typeof row === 'object' && row !== null
364 && (row as { type?: unknown }).type === 'turn/end') throw this.issue
365 return
366 }
367 this.restore.decodeRow(row)
368 }
369}
370
371interface StartedMigrationStream {
372 readonly parser: MigratingJsonlRows
373}
374
375async 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 validation
389 const stream = format.createRestore(header)
390 return { parser: new MigratingJsonlRows(stream) }
391}
392
393async 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_MS
400 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) continue
404 await scheduler.yield()
405 signal?.throwIfAborted()
406 yieldDeadline = performance.now() + MIGRATION_DECODE_YIELD_INTERVAL_MS
407 }
408 }
409}
410
411async 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 }
440
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}
482
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 */
492export 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}
501
502async 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}
552
553/** Fully replay embedded streams only inside isolated current-generation verification. */
554function assertCurrentAssistantStreams(events: readonly SessionEvent[]): void {
555 for (const [index, event] of events.entries()) {
556 if (event.type !== 'assistant/message' && event.type !== 'assistant/attempt') continue
557 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) continue
566 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}
578
579function 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}
607
608function finishCurrentGenerationScan(
609 scanner: SessionLogScanner,
610): ReturnType<SessionLogScanner['finish']> {
611 const inputBytes = scanner.checkpoint().inputBytes
612 const decoded = scanner.finish()
613 if (decoded.committedBytes !== inputBytes) throw new Error('current session generation has a torn physical tail')
614 return decoded
615}
616
617function stringifyJson(value: unknown, subject: string): string {
618 let text: unknown
619 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 text
626}
627
628function 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}
633
634function 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}
654
655async function syncDirectory(path: string, internals: JsonlGenerationInternals): Promise<void> {
656 /* v8 ignore next -- Windows namespace operations request write-through directly. */
657 if (internals.platform === 'win32') return
658 const handle = await internals.fs.open(path, 'r')
659 try {
660 await handle.sync()
661 } finally {
662 await handle.close()
663 }
664}
665
666interface StreamedMigrationStage {
667 readonly path: string
668 readonly bytes: number
669 readonly digest: string
670}
671
672/** Produce bounded JSONL chunks while yielding between main-thread encoding slices. */
673async 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 = 0
681 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 = 0
690 }
691 lines.push(line)
692 bytes += lineBytes
693 }
694 yield Buffer.from(lines.join(''))
695}
696
697async function writeMigrationChunks(
698 chunks: AsyncIterable<Buffer>,
699 write: (chunk: Buffer) => Promise<void>,
700): Promise<void> {
701 let pending: Buffer[] = []
702 let bytes = 0
703 for await (const chunk of chunks) {
704 pending.push(chunk)
705 bytes += chunk.length
706 if (bytes < MIGRATION_WRITE_CHUNK_BYTES) continue
707 await write(pending.length === 1 ? pending[0] as Buffer : Buffer.concat(pending, bytes))
708 pending = []
709 bytes = 0
710 }
711 if (bytes > 0) await write(pending.length === 1 ? pending[0] as Buffer : Buffer.concat(pending, bytes))
712}
713
714/** Encode directly into one synced stage without a whole-artifact row or byte buffer. */
715async 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: string
726 let handle: FileHandle
727 for (;;) {
728 path = join(dirname(currentPath), `session.migration.${internals.randomToken()}${suffix}.tmp`)
729 try {
730 handle = await internals.fs.open(path, 'wx', 0o600)
731 break
732 } catch (error) {
733 if (isEEXIST(error)) continue
734 throw error
735 }
736 }
737 const hash = createHash('sha256')
738 let bytes = 0
739 const write = async (chunk: Buffer): Promise<void> => {
740 await handle.writeFile(chunk)
741 hash.update(chunk)
742 bytes += chunk.length
743 }
744 let failure: unknown
745 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 = error
771 }
772 try {
773 await handle.close()
774 } catch (error: unknown) {
775 failure = failure === undefined
776 ? error
777 : new AggregateError([failure, error], `failed to write and close migration stage "${path}"`)
778 }
779 if (failure !== undefined) {
780 const writeError = failure instanceof Error
781 ? failure
782 : new Error('migration stage write failed with a non-Error rejection', { cause: failure })
783 await removeTemporary(path, writeError, internals)
784 throw writeError
785 }
786 return { path, bytes, digest: hash.digest('hex') }
787}
788
789/** Remove one temporary file without hiding the operation failure that made it disposable. */
790async 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}
804
805/** Remove a redundant stage after the target has been validated as committed. */
806async 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}
816
817async 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 true
826 } catch (error) {
827 /* v8 ignore else -- native helper tests own non-collision Win32 failures. */
828 if (isEEXIST(error)) return false
829 /* v8 ignore next -- the filesystem error is already complete. */
830 throw error
831 }
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 false
838 /* v8 ignore next -- the filesystem error is already complete. */
839 throw error
840 }
841 await syncDirectory(dirname(currentPath), internals)
842 return true
843}
844
845function 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}
850
851async 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 error
872 throw new JsonlGenerationTargetConflictError(currentPath, asError(error))
873 }
874}
875
876function withOverrides(overrides: JsonlGenerationRuntimeOverrides): JsonlGenerationInternals {
877 return {
878 ...defaultInternals,
879 ...overrides,
880 fs: { ...defaultFileSystem, ...overrides.fs },
881 }
882}
883
884async 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 } = options
893 const eventCount = artifact.events.length
894 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: JsonlPhysicalIdentity
915 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 candidate
934 })
935 currentIdentity = winner.identity
936 await removeCommittedTemporary(staged.path, internals)
937 staged = { ...staged, path: '' }
938 }
939 return currentIdentity
940 } catch (error: unknown) {
941 if (staged.path !== '') await removeTemporary(staged.path, error, internals)
942 throw error
943 }
944}
945
946async function prepareMigration(
947 options: PrepareJsonlMigrationOptions,
948 internals: JsonlGenerationInternals,
949): Promise<PreparedJsonlMigration> {
950 const { sourcePath, sourceVersion, currentPath, compression, format, signal } = options
951 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: SessionFormatArtifact
963 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 error
977 }
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.identity
983 let publication: Promise<JsonlPhysicalIdentity> | undefined
984 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 publication
998 },
999 }
1000}
1001
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 */
1007export function prepareJsonlMigration(
1008 options: PrepareJsonlMigrationOptions,
1009): Promise<PreparedJsonlMigration> {
1010 return defaultGenerationRuntime.prepare(options)
1011}
1012
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 */
1018export 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}
1030
1031const defaultGenerationRuntime = createJsonlGenerationRuntime()
1032
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 */
1043export 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: SessionFormatArtifact
1052 try {
1053 artifact = await decodeStreamingMigration(source.bytes, compression, version, format, undefined, signal)
1054 } catch (error: unknown) {
1055 if (signal?.aborted || error instanceof SessionFormatError) throw error
1056 throw new SessionFormatError(String(error), { cause: error })
1057 }
1058 return { artifact, identity: source.identity }
1059}