Răsfoiți Sursa

perf(session-persistence): stream migration publication and verification

imccyu 1 săptămână în urmă
părinte
comite
ec2f63dbdb

+ 2 - 0
packages/session/session-persistence-jsonl/package.json

@@ -23,6 +23,7 @@
   },
   "files": [
     "lib/index.js",
+    "lib/worker.cjs",
     "lib/types/**/*.d.ts"
   ],
   "license": "MIT",
@@ -32,6 +33,7 @@
     "@deepseek-ai/cordis": "workspace:^"
   },
   "dependencies": {
+    "@deepseek-ai/dsh-llm": "workspace:^",
     "@deepseek-ai/dsh-session-format": "workspace:^",
     "@deepseek-ai/dsh-session-format-catalog": "workspace:^",
     "fs-ext": "2.1.1",

+ 61 - 82
packages/session/session-persistence-jsonl/src/format.ts

@@ -10,7 +10,7 @@
 
 import { isAbsolute, join } from 'node:path'
 import {
-  decodeSeqRanges, encodeSeqRanges, SESSION_FORMAT_VERSION,
+  SESSION_FORMAT_VERSION,
   SessionLogOffset,
 } from '@deepseek-ai/dsh-session'
 import type {
@@ -20,6 +20,9 @@ import type {
   SessionLogOffset as SessionLogOffsetType,
 } from '@deepseek-ai/dsh-session'
 import { parseSessionFormatLogFilename, sessionFormatLogFilename } from '@deepseek-ai/dsh-session-format'
+import type { SessionFormatEvent } from '@deepseek-ai/dsh-session-format'
+import type { SessionFormatRecovery, SessionFormatRestore } from '@deepseek-ai/dsh-session-format'
+import { sessionFormatCatalog } from '@deepseek-ai/dsh-session-format-catalog'
 import {
   SessionFormatUnsupportedError,
   sessionFormatVersionRefusal,
@@ -122,18 +125,10 @@ export function toHeaderLine(
   if (!header.isSeeded && cut !== 0) {
     throw new Error('unseeded session header inherited event count must be 0')
   }
-  return {
-    type: 'session',
-    version: header.version,
-    id: header.id,
-    createdAt: header.createdAt,
-    ...header.cwd !== undefined ? { cwd: header.cwd } : {},
-    ...header.parentSession !== undefined ? { parentSession: header.parentSession } : {},
-    isSeeded: header.isSeeded,
-    ...header.origin !== undefined ? { origin: header.origin } : {},
+  return sessionFormatCatalog.encodeCurrentHeader({
+    ...header,
     delegationDepth: header.delegationDepth ?? 0,
-    ...header.agentPreset !== undefined ? { agentPreset: header.agentPreset } : {},
-  }
+  }, cut) as unknown as HeaderLine
 }
 
 /**
@@ -314,38 +309,16 @@ export function logPath(
  * @returns the batch's JSONL text; the writer adds the final newline.
  */
 export function eventLines(events: readonly SessionEvent[]): string {
-  return events.map(record => JSON.stringify(encodeProvenanceForStorage(record))).join('\n')
-}
-
-/**
- * Losslessly shrink a record's `sourceEventSeqs` for the log: consecutive
- * runs of at least three seqs become `[start, end]` pairs, and any other list
- * stays verbatim.
- * @param record - one stored record (event or packed row).
- * @returns the record with its provenance in storage form (widened from the
- *   in-memory `SessionSeq[]`; {@link expandProvenanceFromStorage} restores it).
- */
-function encodeProvenanceForStorage(record: SessionEvent): unknown {
-  if (!('sourceEventSeqs' in record)) return record
-  return { ...record, sourceEventSeqs: encodeSeqRanges(record.sourceEventSeqs) }
+  return events.map(eventLine).join('\n')
 }
 
 /**
- * Expand a parsed line's storage-form provenance back to `SessionSeq[]`.
- * @param parsed - the JSON-parsed value of one stored line.
- * @returns the value with provenance expanded.
- * @throws when the record or its storage-form provenance is malformed.
+ * Serialize one v2 event as one JSONL record without its trailing newline.
+ * @param event - current event to encode.
+ * @returns one physical JSON record.
  */
-function expandProvenanceFromStorage(parsed: unknown): unknown {
-  if (typeof parsed !== 'object' || parsed === null || Array.isArray(parsed)) {
-    throw new TypeError('stored session records must be objects')
-  }
-  const record = parsed as { seq?: unknown; sourceEventSeqs?: unknown }
-  if (record.sourceEventSeqs === undefined) return parsed
-  if (!Number.isSafeInteger(record.seq) || (record.seq as number) < 0) {
-    throw new TypeError('stored session event seq must be a non-negative safe integer')
-  }
-  return { ...record, sourceEventSeqs: decodeSeqRanges(record.sourceEventSeqs, record.seq as number) }
+export function eventLine(event: SessionEvent): string {
+  return JSON.stringify(sessionFormatCatalog.encodeCurrentEvent(event as unknown as SessionFormatEvent))
 }
 
 interface SessionLogScan {
@@ -355,21 +328,6 @@ interface SessionLogScan {
   committedBytes: number
 }
 
-/** Derive the v2 fork cut from the last lineage-tagged seed marker. */
-function inheritedCut(meta: SessionHeader, events: readonly SessionEvent[]): SessionLogOffsetType {
-  let cut: SessionLogOffsetType | undefined
-  for (const event of events) {
-    if (event.type === 'session/end-seed' && event.data.inherited === true) cut = SessionLogOffset(event.seq)
-  }
-  if (meta.isSeeded && cut === undefined) {
-    throw new Error('corrupt session log: seeded v2 header lacks an inherited end-seed marker')
-  }
-  if (!meta.isSeeded && cut !== undefined) {
-    throw new Error('corrupt session log: unseeded v2 header contains an inherited end-seed marker')
-  }
-  return cut ?? SessionLogOffset(0)
-}
-
 /**
  * Refuse a header carrying a format version this build does not read BEFORE
  * validating the current header shape or decoding any event row: a future
@@ -377,8 +335,7 @@ function inheritedCut(meta: SessionHeader, events: readonly SessionEvent[]): Ses
  * must see "upgrade the harness", never "corrupt session log".
  * @param parsed - the JSON-parsed first line of a session artifact.
  */
-function refuseForeignFormatVersion(parsed: unknown): void {
-  if (typeof parsed !== 'object' || parsed === null) return
+function refuseForeignFormatVersion(parsed: object): void {
   const { version, id } = parsed as { version?: unknown; id?: unknown }
   if (typeof version !== 'number' || version === SESSION_FORMAT_VERSION) return
   throw new SessionFormatUnsupportedError(
@@ -387,7 +344,7 @@ function refuseForeignFormatVersion(parsed: unknown): void {
 }
 
 /** Parse one complete header record supplied independently from event rows. */
-function parseHeaderRecord(record: Buffer): ReturnType<typeof fromHeaderLine> {
+function parseHeaderRecord(record: Buffer): { readonly meta: SessionHeader; readonly restore: SessionFormatRestore } {
   if (record.length === 0 || record.at(-1) !== 0x0A || record.indexOf(0x0A) !== record.length - 1) {
     throw new Error('empty or header-less session log')
   }
@@ -397,12 +354,25 @@ function parseHeaderRecord(record: Buffer): ReturnType<typeof fromHeaderLine> {
   } catch {
     throw new Error('corrupt session log: header line is not valid JSON')
   }
+  if (typeof parsed !== 'object' || parsed === null || Array.isArray(parsed)) {
+    throw new Error('corrupt session log: first line is not a JSON object')
+  }
   refuseForeignFormatVersion(parsed)
   assertNoRetiredHeaderFields(parsed)
   if (!isHeaderLine(parsed)) {
     throw new Error('corrupt session log: first line is not a session header')
   }
-  return fromHeaderLine(parsed)
+  let restore: SessionFormatRestore
+  try {
+    restore = sessionFormatCatalog.createRestore(parsed, {
+      recovery: 'strict',
+      validation: 'transformed',
+    })
+  } catch {
+    /* v8 ignore next -- isHeaderLine matches the current codec; this preserves classification if it tightens. */
+    throw new Error('corrupt session log: first line is not a session header')
+  }
+  return { meta: fromHeaderLine(parsed).meta, restore }
 }
 
 /**
@@ -413,7 +383,8 @@ function parseHeaderRecord(record: Buffer): ReturnType<typeof fromHeaderLine> {
  */
 export class SessionLogScanner {
   private readonly meta: SessionHeader
-  private readonly events: SessionEvent[] = []
+  private readonly restore: SessionFormatRestore
+  private eventCount = 0
   private fragments: Buffer[] = []
   private fragmentBytes = 0
   private inputBytes: number
@@ -426,9 +397,13 @@ export class SessionLogScanner {
    * Create an event scanner from exactly one newline-terminated header record.
    * @param headerRecord - the complete first JSONL record, including its newline.
    */
-  constructor(headerRecord: Buffer) {
+  constructor(
+    headerRecord: Buffer,
+    private readonly recovery: SessionFormatRecovery = 'recoverable',
+  ) {
     const parsed = parseHeaderRecord(headerRecord)
     this.meta = parsed.meta
+    this.restore = parsed.restore
     this.inputBytes = headerRecord.length
     this.committedBytes = headerRecord.length
   }
@@ -477,7 +452,7 @@ export class SessionLogScanner {
     return {
       inputBytes: this.inputBytes,
       committedBytes: this.committedBytes,
-      eventCount: SessionLogOffset(this.events.length),
+      eventCount: SessionLogOffset(this.eventCount),
     }
   }
 
@@ -487,10 +462,11 @@ export class SessionLogScanner {
    */
   finish(): SessionLogScan {
     this.finished = true
+    const artifact = this.restore.finish()
     return {
       meta: this.meta,
-      inheritedEventCount: inheritedCut(this.meta, this.events),
-      events: this.events,
+      inheritedEventCount: SessionLogOffset(artifact.inheritedEventCount),
+      events: artifact.events as unknown as SessionEvent[],
       committedBytes: this.committedBytes,
     }
   }
@@ -498,33 +474,36 @@ export class SessionLogScanner {
   /** Decode one complete event row and update the contiguous prefix. */
   private consumeEventLine(line: Buffer, endByte: number): void {
     this.eventLine += 1
-    let decoded: SessionEvent[]
+    let decoded: unknown
     try {
-      decoded = [expandProvenanceFromStorage(JSON.parse(line.toString('utf8'))) as SessionEvent]
+      decoded = JSON.parse(line.toString('utf8')) as unknown
     } catch {
-      this.issue ??= new Error(`corrupt session log: unparsable committed event at line ${this.eventLine}`)
+      const issue = new Error(`corrupt session log: unparsable committed event at line ${this.eventLine}`)
+      if (this.recovery === 'strict') throw issue
+      this.issue ??= issue
       return
     }
 
     if (this.issue !== undefined) {
-      if (decoded.some(event => event.type === 'turn/end')) throw this.issue
+      if (typeof decoded === 'object' && decoded !== null
+        && (decoded as { type?: unknown }).type === 'turn/end') throw this.issue
       return
     }
-
-    const rowStart = this.events.length
-    for (const event of decoded) {
-      if (event.seq !== this.events.length) {
-        const expected = this.events.length
-        this.events.length = rowStart
-        this.issue = new Error(
-          `corrupt session log: seq gap in committed region at line ${this.eventLine} `
-          + `(expected ${expected}, got ${event.seq})`,
-        )
-        if (decoded.some(candidate => candidate.type === 'turn/end')) throw this.issue
-        return
-      }
-      this.events.push(event)
+    try {
+      this.restore.decodeRow(decoded)
+    } catch (error: unknown) {
+      /* v8 ignore next -- every production Session format decoder rejects with Error. */
+      const detail = error instanceof Error ? error.message : String(error)
+      const issue = new Error(`corrupt session log: invalid committed event at line ${this.eventLine}: ${detail}`, {
+        cause: error,
+      })
+      if (this.recovery === 'strict') throw issue
+      this.issue = issue
+      if (typeof decoded === 'object' && decoded !== null
+        && (decoded as { type?: unknown }).type === 'turn/end') throw issue
+      return
     }
+    this.eventCount += 1
     this.committedBytes = endByte
   }
 }

+ 541 - 229
packages/session/session-persistence-jsonl/src/generation.ts

@@ -19,33 +19,45 @@ import {
   type FileHandle,
 } from 'node:fs/promises'
 import { basename, dirname, join } from 'node:path'
+import { performance } from 'node:perf_hooks'
+import { pipeline, Readable } from 'node:stream'
+import { scheduler } from 'node:timers/promises'
+import { constants, createZstdCompress } from 'node:zlib'
+import { Session } from '@deepseek-ai/dsh-session'
+import type {
+  SessionFormatArtifact,
+  SessionFormatJsonValue,
+  SessionFormatRestore,
+} from '@deepseek-ai/dsh-session-format'
+import { validateStoredEvents } from '@deepseek-ai/dsh-session-persistence'
 import type { JsonlCompression } from './format.ts'
-import { generationLogFilename, logSuffix } from './format.ts'
+import { generationLogFilename, logSuffix, SessionLogScanner } from './format.ts'
 import { publishNewFileWin32 } from './win32.ts'
 import {
   compressZstdFrame,
   createZstdFrameDecoder,
-  decompressZstdFrame,
   decompressZstdPrefix,
   scanZstdFrames,
 } from './zstd.ts'
 
-/** Parsed JSONL values supplied to the format catalog. */
-export interface JsonlDecodedGeneration {
-  readonly header: Record<string, unknown>
-  readonly rows: readonly unknown[]
+/** Internal scheduling bounds: preserve old decode cadence and cap each synchronous encode slice. */
+const MIGRATION_DECODE_YIELD_INTERVAL_MS = 500
+const MIGRATION_WORK_CHUNK_BYTES = 1024 * 1024
+const MIGRATION_WRITE_CHUNK_BYTES = 4 * 1024 * 1024
+const ZSTD_CHECKSUM_OPTIONS = {
+  chunkSize: MIGRATION_WORK_CHUNK_BYTES,
+  params: { [constants.ZSTD_c_checksumFlag]: 1 },
 }
 
-/** Current JSONL values returned by the format catalog for physical encoding. */
-export interface JsonlCurrentGeneration extends JsonlDecodedGeneration {}
-
 /** Pure adapter between backend-owned JSONL framing and the format catalog. */
 export interface JsonlGenerationFormatAdapter {
   readonly currentVersion: number
-  /** Convert one detached historical generation to exact current JSON values. */
-  migrate(source: JsonlDecodedGeneration): JsonlCurrentGeneration
-  /** Validate one decoded current generation, including after committed reopen. */
-  validateCurrent(candidate: JsonlCurrentGeneration): void
+  /** Create the single-pass codec and migration state for a historical header. */
+  createRestore(header: Record<string, unknown>): SessionFormatRestore
+  /** Encode one current header record without materializing body rows. */
+  encodeHeader(header: SessionFormatArtifact['header'], inheritedEventCount: number): SessionFormatJsonValue
+  /** Encode one current event record. */
+  encodeEvent(event: SessionFormatArtifact['events'][number]): SessionFormatJsonValue
   /** Classify a supported-version artifact that policy refuses to migrate. */
   isUnsupportedMigrationError?(error: unknown): error is Error
 }
@@ -64,9 +76,31 @@ export interface EnsureJsonlGenerationOptions {
   readonly validateHistoricalHeader?: (
     header: Readonly<Record<string, unknown>>,
   ) => void | Promise<void>
+  /** Validate the staged file in an isolated worker before publication. */
+  readonly verifyCurrentFile: (
+    path: string,
+    compression: JsonlCompression,
+    expectedId: string,
+    expectedEventCount: number,
+    expectedPrefix?: JsonlExpectedPrefix,
+    signal?: AbortSignal,
+  ) => Promise<JsonlVerifiedGeneration>
   readonly signal?: AbortSignal
 }
 
+/** Small physical identity returned by an isolated generation verifier. */
+export interface JsonlVerifiedGeneration {
+  readonly identity: JsonlPhysicalIdentity
+  readonly bytes: number
+  readonly digest: string
+}
+
+/** Physical byte prefix already proven to be a valid complete generation. */
+export interface JsonlExpectedPrefix {
+  readonly bytes: number
+  readonly digest: string
+}
+
 /** Result of current classification or exclusive publication. */
 export type EnsureJsonlGenerationResult =
   | {
@@ -88,11 +122,6 @@ export type EnsureJsonlGenerationResult =
 export class JsonlGenerationNewerVersionError extends Error {
   override readonly name = 'JsonlGenerationNewerVersionError'
 
-  /**
-   * @param storedVersion - version read from the highest stored generation.
-   * @param currentVersion - version this build writes.
-   * @param storedId - minimally decoded identity used in the refusal diagnostic.
-   */
   constructor(
     readonly storedVersion: number,
     readonly currentVersion: number,
@@ -143,10 +172,8 @@ export interface JsonlPhysicalIdentity {
   readonly ctimeNs: bigint
 }
 
-/** One revision-stable physical artifact read reusable by the immediate backend hook. */
-export interface JsonlPhysicalSnapshot {
-  readonly bytes: Buffer
-  readonly identity: JsonlPhysicalIdentity
+/** One revision-stable physical artifact returned to the immediate backend decoder. */
+export interface JsonlPhysicalSnapshot extends StablePhysicalFile {
   readonly headerValue: Record<string, unknown>
   readonly headerRecord: Buffer
 }
@@ -162,11 +189,6 @@ interface JsonlPhysicalHeader {
   readonly record: Buffer
 }
 
-interface DecodedPhysicalJsonl {
-  readonly bytes: Buffer
-  readonly torn: boolean
-}
-
 interface GenerationFileSystem {
   open(path: string, flags: string, mode?: number): Promise<FileHandle>
   readFile(path: string, signal?: AbortSignal): Promise<Buffer>
@@ -189,10 +211,24 @@ interface JsonlGenerationInternals {
   readonly barrier: (phase: GenerationBarrierPhase, attempt: number) => void | Promise<void>
 }
 
-type JsonlGenerationTestOverrides = Partial<Omit<JsonlGenerationInternals, 'fs'>> & {
+/** Dependency overrides for an isolated generation runtime. */
+export type JsonlGenerationRuntimeOverrides = Partial<Omit<JsonlGenerationInternals, 'fs'>> & {
   readonly fs?: Partial<GenerationFileSystem>
 }
 
+/** Bound generation operations used by production defaults and deterministic tests. */
+export interface JsonlGenerationRuntime {
+  readStable(path: string, signal?: AbortSignal): Promise<StablePhysicalFile>
+  ensure(options: EnsureJsonlGenerationOptions): Promise<EnsureJsonlGenerationResult>
+  verify(
+    path: string,
+    compression: JsonlCompression,
+    expectedId: string,
+    expectedEventCount: number,
+    expectedPrefix?: JsonlExpectedPrefix,
+  ): Promise<JsonlVerifiedGeneration>
+}
+
 const defaultFileSystem: GenerationFileSystem = {
   open: (path, flags, mode) => fsOpen(path, flags, mode),
   readFile: (path, signal) => fsReadFile(path, signal === undefined ? undefined : { signal }),
@@ -240,7 +276,7 @@ export async function readStableJsonlFile(
   path: string,
   signal?: AbortSignal,
 ): Promise<StablePhysicalFile> {
-  return readStableSnapshot(path, signal, defaultFileSystem)
+  return defaultGenerationRuntime.readStable(path, signal)
 }
 
 async function readStableSnapshot(
@@ -289,120 +325,299 @@ function parseJson(text: string, subject: string): unknown {
   }
 }
 
-function parseGeneration(bytes: Buffer, recoverSuffix = false): JsonlDecodedGeneration {
-  /* v8 ignore next -- decodePhysicalJsonl supplies a non-empty newline-terminated prefix. */
-  if (bytes.length === 0 || bytes.at(-1) !== 0x0A) {
-    throw new Error('empty or header-less session log')
+/** Incremental JSONL parser that retains only one cross-frame record fragment. */
+class MigratingJsonlRows {
+  private fragments: Buffer[] = []
+  private fragmentBytes = 0
+  private rowIndex = 0
+  private issue: Error | undefined
+
+  constructor(private readonly restore: SessionFormatRestore) {}
+
+  /** Consume plaintext bytes following the independently decoded header. */
+  write(chunk: Buffer): void {
+    /* jscpd:ignore-start -- migration parsing and readable-log scanning own different recovery and byte-accounting state. */
+    let lineStart = 0
+    for (
+      let newline = chunk.indexOf(0x0A);
+      newline !== -1;
+      newline = chunk.indexOf(0x0A, lineStart)
+    ) {
+      const fragment = chunk.subarray(lineStart, newline)
+      let line = fragment
+      if (this.fragments.length > 0) {
+        if (fragment.length > 0) this.fragments.push(fragment)
+        line = Buffer.concat(this.fragments, this.fragmentBytes + fragment.length)
+        this.fragments = []
+        this.fragmentBytes = 0
+      }
+      this.consume(line)
+      lineStart = newline + 1
+    }
+    if (lineStart < chunk.length) {
+      const fragment = Buffer.from(chunk.subarray(lineStart))
+      this.fragments.push(fragment)
+      this.fragmentBytes += fragment.length
+    }
+    /* jscpd:ignore-end */
   }
-  const records = bytes.toString('utf8').slice(0, -1).split('\n')
-  const parsedHeader = parseJson(records[0] as string, 'header line')
-  storedVersion(parsedHeader)
-  const rows: unknown[] = []
-  let issue: Error | undefined
-  for (const [index, record] of records.slice(1).entries()) {
+
+  /** Refuse a record fragment left by structurally complete Zstandard frames. */
+  assertCompleteFramesEndOnRecord(): void {
+    if (this.fragments.length > 0) {
+      throw new Error('corrupt Zstandard session log: complete frame contains a torn JSONL record')
+    }
+  }
+
+  finish(): SessionFormatArtifact {
+    return this.restore.finish()
+  }
+
+  private consume(line: Buffer): void {
+    const index = this.rowIndex
+    this.rowIndex += 1
     let row: unknown
     try {
-      row = parseJson(record, `row ${index + 1}`)
-    } catch (error) {
-      if (!recoverSuffix) throw error
-      issue ??= error as Error
-      continue
+      row = parseJson(line.toString('utf8'), `row ${index + 1}`)
+    } catch (error: unknown) {
+      this.issue ??= asError(error)
+      return
     }
-    if (issue !== undefined) {
+    if (this.issue !== undefined) {
       if (typeof row === 'object' && row !== null
-        && (row as { type?: unknown }).type === 'turn/end') throw issue
-      continue
+        && (row as { type?: unknown }).type === 'turn/end') throw this.issue
+      return
     }
-    rows.push(row)
+    this.restore.decodeRow(row)
   }
-  return { header: parsedHeader as Record<string, unknown>, rows }
 }
 
-function stringifyJson(value: unknown, subject: string): string {
-  let text: unknown
-  try {
-    text = JSON.stringify(value)
-  } catch (error) {
-    throw new Error(`${subject} is not lossless JSON`, { cause: error })
-  }
-  if (typeof text !== 'string') throw new Error(`${subject} is not lossless JSON`)
-  return text
+interface StartedMigrationStream {
+  readonly parser: MigratingJsonlRows
 }
 
-function encodeLogicalJsonl(generation: JsonlCurrentGeneration): Buffer {
-  const records = [
-    stringifyJson(generation.header, 'migrated session header'),
-    ...generation.rows.map((row, index) => stringifyJson(row, `migrated session row ${index + 1}`)),
-  ]
-  return Buffer.from(`${records.join('\n')}\n`)
+async function startMigrationStream(
+  headerRecord: Buffer,
+  format: JsonlGenerationFormatAdapter,
+  validateHistoricalHeader?: EnsureJsonlGenerationOptions['validateHistoricalHeader'],
+): Promise<StartedMigrationStream> {
+  const value = parseJson(headerRecord.subarray(0, -1).toString('utf8'), 'header line')
+  const header = value as Record<string, unknown>
+  const validation = validateHistoricalHeader?.(header)
+  if (validation !== undefined) await validation
+  const stream = format.createRestore(header)
+  return { parser: new MigratingJsonlRows(stream) }
 }
 
-function assertIndependentHeaderFrame(plaintext: Buffer): void {
-  if (plaintext.length === 0 || plaintext.indexOf(0x0A) !== plaintext.length - 1) {
-    throw new Error('corrupt Zstandard session log: first frame is not exactly one header line')
+async function consumeMigrationBytes(
+  rows: MigratingJsonlRows,
+  chunks: Iterable<Buffer>,
+  signal?: AbortSignal,
+): Promise<void> {
+  signal?.throwIfAborted()
+  let yieldDeadline = performance.now() + MIGRATION_DECODE_YIELD_INTERVAL_MS
+  for (const bytes of chunks) {
+    for (let offset = 0; offset < bytes.length; offset += MIGRATION_WORK_CHUNK_BYTES) {
+      rows.write(bytes.subarray(offset, offset + MIGRATION_WORK_CHUNK_BYTES))
+      if (performance.now() < yieldDeadline) continue
+      await scheduler.yield()
+      signal?.throwIfAborted()
+      yieldDeadline = performance.now() + MIGRATION_DECODE_YIELD_INTERVAL_MS
+    }
   }
 }
 
-async function decodeZstdJsonl(bytes: Buffer, signal?: AbortSignal): Promise<DecodedPhysicalJsonl> {
+async function decodeStreamingMigration(
+  bytes: Buffer,
+  compression: JsonlCompression,
+  format: JsonlGenerationFormatAdapter,
+  validateHistoricalHeader: EnsureJsonlGenerationOptions['validateHistoricalHeader'],
+  signal?: AbortSignal,
+): Promise<SessionFormatArtifact> {
   signal?.throwIfAborted()
+  if (compression === 'none') {
+    const headerEnd = bytes.indexOf(0x0A)
+    /* v8 ignore next -- ensureCurrent's physical-header preflight already requires this newline. */
+    if (headerEnd === -1) throw new Error('empty or header-less session log')
+    const stream = await startMigrationStream(
+      bytes.subarray(0, headerEnd + 1),
+      format,
+      validateHistoricalHeader,
+    )
+    signal?.throwIfAborted()
+    const bodyEnd = bytes.lastIndexOf(0x0A)
+    if (bodyEnd > headerEnd) {
+      await consumeMigrationBytes(
+        stream.parser,
+        [bytes.subarray(headerEnd + 1, bodyEnd + 1)],
+        signal,
+      )
+    }
+    return stream.parser.finish()
+  }
+
   const { frames, tornStart } = scanZstdFrames(bytes)
-  /* v8 ignore next -- the independent header probe already established the first frame. */
+  /* v8 ignore next -- ensureCurrent's physical-header preflight already requires a complete header frame. */
   if (frames.length === 0) throw new Error('empty or header-less Zstandard session log')
-  const complete: Buffer[] = []
-  for (const [index, frame] of frames.entries()) {
+  const decoder = createZstdFrameDecoder()
+  try {
+    const decoded = decoder.decode(bytes, frames)
+    const first = decoded.next()
+    /* v8 ignore next -- a non-empty structural frame list yields once or throws. */
+    if (first.done) throw new Error('empty or header-less Zstandard session log')
+    assertIndependentHeaderFrame(first.value)
+    const stream = await startMigrationStream(
+      first.value,
+      format,
+      validateHistoricalHeader,
+    )
     signal?.throwIfAborted()
-    const plaintext = await decompressZstdFrame(bytes.subarray(frame.start, frame.end))
-    if (index === 0) assertIndependentHeaderFrame(plaintext)
-    complete.push(plaintext)
-  }
-  const completeBytes = Buffer.concat(complete)
-  if (completeBytes.at(-1) !== 0x0A) {
-    throw new Error('corrupt Zstandard session log: complete frame contains a torn JSONL record')
+    await consumeMigrationBytes(stream.parser, decoded, signal)
+    stream.parser.assertCompleteFramesEndOnRecord()
+    if (tornStart !== undefined) {
+      let recovered: Buffer = Buffer.alloc(0)
+      try {
+        recovered = await decompressZstdPrefix(bytes.subarray(tornStart))
+      } catch {
+        /* v8 ignore next -- decoder failure plus concurrent abort is timing-dependent. */
+        if (signal?.aborted) signal.throwIfAborted()
+      }
+      signal?.throwIfAborted()
+      const newline = recovered.lastIndexOf(0x0A)
+      if (newline !== -1) {
+        await consumeMigrationBytes(
+          stream.parser,
+          [recovered.subarray(0, newline + 1)],
+          signal,
+        )
+      }
+    }
+    return stream.parser.finish()
+  } finally {
+    decoder.close()
   }
-  if (tornStart === undefined) return { bytes: completeBytes, torn: false }
+}
 
-  let recovered = Buffer.alloc(0)
-  try {
-    recovered = Buffer.from(await decompressZstdPrefix(bytes.subarray(tornStart)))
-  } catch {
-    /* v8 ignore next -- an abort racing decoder failure is timing-dependent */
-    if (signal?.aborted) signal.throwIfAborted()
-    // A structurally torn frame may produce no plaintext; prior frames remain valid.
+/**
+ * Read and validate one complete current generation for an isolated verifier.
+ * @param path - staged or competing current-generation path.
+ * @param compression - configured physical encoding.
+ * @param expectedId - Session identity expected in the header.
+ * @param expectedEventCount - exact logical event count expected after decoding.
+ * @param expectedPrefix - verified migration prefix; an append tail may be present and is not validated.
+ * @returns stable physical identity and digest for publication comparison.
+ */
+export async function verifyJsonlCurrentGeneration(
+  path: string,
+  compression: JsonlCompression,
+  expectedId: string,
+  expectedEventCount: number,
+  expectedPrefix?: JsonlExpectedPrefix,
+): Promise<JsonlVerifiedGeneration> {
+  return defaultGenerationRuntime.verify(path, compression, expectedId, expectedEventCount, expectedPrefix)
+}
+
+async function verifyCurrentGeneration(
+  path: string,
+  compression: JsonlCompression,
+  expectedId: string,
+  expectedEventCount: number,
+  fs: GenerationFileSystem,
+  expectedPrefix?: JsonlExpectedPrefix,
+): Promise<JsonlVerifiedGeneration> {
+  const before = await fs.stat(path)
+  const bytes = await fs.readFile(path)
+  const after = await fs.stat(path)
+  if (expectedPrefix !== undefined) {
+    if (bytes.length < expectedPrefix.bytes) {
+      throw new Error('target bytes are shorter than the migrated generation')
+    }
+    const digest = createHash('sha256').update(bytes.subarray(0, expectedPrefix.bytes)).digest('hex')
+    if (digest !== expectedPrefix.digest) {
+      throw new Error('target bytes do not begin with the migrated generation')
+    }
+    return { identity: after, bytes: expectedPrefix.bytes, digest }
   }
-  signal?.throwIfAborted()
-  const newline = recovered.lastIndexOf(0x0A)
+  if (identity(before) !== identity(after)) {
+    throw new Error('current session generation changed during verification')
+  }
+  const snapshot = { bytes, identity: after }
+  const generation = decodeCurrentGeneration(snapshot.bytes, compression)
+  validateStoredEvents(generation.meta, generation.events, { kind: 'jsonl', path })
+  if (generation.meta.id !== expectedId) {
+    throw new Error(`current session generation contains id "${generation.meta.id}", expected "${expectedId}"`)
+  }
+  if (generation.events.length !== expectedEventCount) {
+    throw new Error(
+      `current session generation contains ${generation.events.length} events, expected ${expectedEventCount}`,
+    )
+  }
+  Session.fromRestore(
+    generation.meta.id,
+    generation.events,
+    generation.meta,
+    generation.inheritedEventCount,
+  )
   return {
-    bytes: newline === -1
-      ? completeBytes
-      : Buffer.concat([completeBytes, recovered.subarray(0, newline + 1)]),
-    torn: true,
+    identity: snapshot.identity,
+    bytes: snapshot.bytes.length,
+    digest: createHash('sha256').update(snapshot.bytes).digest('hex'),
   }
 }
 
-async function decodePhysicalJsonl(
+function decodeCurrentGeneration(
   bytes: Buffer,
   compression: JsonlCompression,
-  signal?: AbortSignal,
-): Promise<DecodedPhysicalJsonl> {
-  if (compression === 'zstd') return decodeZstdJsonl(bytes, signal)
-  signal?.throwIfAborted()
-  const newline = bytes.lastIndexOf(0x0A)
-  /* v8 ignore next -- physical header classification already found a newline in the same stable bytes. */
-  if (newline === -1) throw new Error('empty or header-less session log')
-  return { bytes: bytes.subarray(0, newline + 1), torn: newline + 1 !== bytes.length }
+): ReturnType<SessionLogScanner['finish']> {
+  if (compression === 'none') {
+    const headerEnd = bytes.indexOf(0x0A)
+    if (headerEnd === -1) throw new Error('empty or header-less session log')
+    const scanner = new SessionLogScanner(bytes.subarray(0, headerEnd + 1), 'strict')
+    scanner.write(bytes.subarray(headerEnd + 1))
+    return finishCurrentGenerationScan(scanner)
+  }
+  const { frames, tornStart } = scanZstdFrames(bytes)
+  if (frames.length === 0) throw new Error('empty or header-less Zstandard session log')
+  if (tornStart !== undefined) throw new Error('current session generation has a torn physical tail')
+  const decoder = createZstdFrameDecoder()
+  try {
+    const plaintext = decoder.decode(bytes, frames)
+    const header = plaintext.next()
+    /* v8 ignore next -- a non-empty structural frame list yields once or throws. */
+    if (header.done) throw new Error('empty or header-less Zstandard session log')
+    assertIndependentHeaderFrame(header.value)
+    const scanner = new SessionLogScanner(header.value, 'strict')
+    for (const chunk of plaintext) scanner.write(chunk)
+    return finishCurrentGenerationScan(scanner)
+  } finally {
+    decoder.close()
+  }
 }
 
-async function encodePhysicalJsonl(
-  logical: Buffer,
-  generation: JsonlCurrentGeneration,
-  compression: JsonlCompression,
-): Promise<Buffer> {
-  if (compression === 'none') return logical
-  const header = Buffer.from(`${stringifyJson(generation.header, 'migrated session header')}\n`)
-  const headerFrame = await compressZstdFrame(header)
-  if (generation.rows.length === 0) return headerFrame
-  const body = logical.subarray(header.length)
-  return Buffer.concat([headerFrame, await compressZstdFrame(body)])
+function finishCurrentGenerationScan(
+  scanner: SessionLogScanner,
+): ReturnType<SessionLogScanner['finish']> {
+  const inputBytes = scanner.checkpoint().inputBytes
+  const decoded = scanner.finish()
+  if (decoded.committedBytes !== inputBytes) throw new Error('current session generation has a torn physical tail')
+  return decoded
+}
+
+function stringifyJson(value: unknown, subject: string): string {
+  let text: unknown
+  try {
+    text = JSON.stringify(value)
+  } catch (error) {
+    throw new Error(`${subject} is not lossless JSON`, { cause: error })
+  }
+  if (typeof text !== 'string') throw new Error(`${subject} is not lossless JSON`)
+  return text
+}
+
+function assertIndependentHeaderFrame(plaintext: Buffer): void {
+  if (plaintext.length === 0 || plaintext.indexOf(0x0A) !== plaintext.length - 1) {
+    throw new Error('corrupt Zstandard session log: first frame is not exactly one header line')
+  }
 }
 
 function readRawHeader(bytes: Buffer): JsonlPhysicalHeader {
@@ -441,8 +656,7 @@ function readPhysicalHeader(
   compression: JsonlCompression,
   signal: AbortSignal | undefined,
 ): JsonlPhysicalHeader {
-  if (compression === 'zstd') return readZstdHeader(bytes, signal)
-  return readRawHeader(bytes)
+  return compression === 'zstd' ? readZstdHeader(bytes, signal) : readRawHeader(bytes)
 }
 
 function assertGenerationPaths(
@@ -477,48 +691,127 @@ async function syncDirectory(path: string, internals: JsonlGenerationInternals):
   }
 }
 
+interface StreamedMigrationStage {
+  readonly path: string
+  readonly bytes: number
+  readonly digest: string
+}
+
+/** Produce bounded JSONL chunks while yielding between main-thread encoding slices. */
+async function* encodeMigrationRows(
+  artifact: SessionFormatArtifact,
+  format: JsonlGenerationFormatAdapter,
+  signal?: AbortSignal,
+): AsyncGenerator<Buffer, void, void> {
+  signal?.throwIfAborted()
+  let lines: string[] = []
+  let bytes = 0
+  for (const value of artifact.events) {
+    const line = `${stringifyJson(format.encodeEvent(value), `migrated Session event ${value.seq}`)}\n`
+    const lineBytes = Buffer.byteLength(line)
+    if (bytes > 0 && bytes + lineBytes > MIGRATION_WORK_CHUNK_BYTES) {
+      yield Buffer.from(lines.join(''))
+      await scheduler.yield()
+      signal?.throwIfAborted()
+      lines = []
+      bytes = 0
+    }
+    lines.push(line)
+    bytes += lineBytes
+  }
+  yield Buffer.from(lines.join(''))
+}
+
+async function writeMigrationChunks(
+  chunks: AsyncIterable<Buffer>,
+  write: (chunk: Buffer) => Promise<void>,
+): Promise<void> {
+  let pending: Buffer[] = []
+  let bytes = 0
+  for await (const chunk of chunks) {
+    pending.push(chunk)
+    bytes += chunk.length
+    if (bytes < MIGRATION_WRITE_CHUNK_BYTES) continue
+    await write(pending.length === 1 ? pending[0] as Buffer : Buffer.concat(pending, bytes))
+    pending = []
+    bytes = 0
+  }
+  if (bytes > 0) await write(pending.length === 1 ? pending[0] as Buffer : Buffer.concat(pending, bytes))
+}
+
+/** Encode directly into one synced stage without a whole-artifact row or byte buffer. */
 async function writeSyncedTemp(
   currentPath: string,
   suffix: string,
-  bytes: Buffer,
+  compression: JsonlCompression,
+  artifact: SessionFormatArtifact,
+  format: JsonlGenerationFormatAdapter,
+  signal: AbortSignal | undefined,
   internals: JsonlGenerationInternals,
-): Promise<string> {
+): Promise<StreamedMigrationStage> {
+  signal?.throwIfAborted()
+  let path: string
+  let handle: FileHandle
   for (;;) {
-    const path = join(dirname(currentPath), `session.migration.${internals.randomToken()}${suffix}.tmp`)
-    let handle: FileHandle
+    path = join(dirname(currentPath), `session.migration.${internals.randomToken()}${suffix}.tmp`)
     try {
       handle = await internals.fs.open(path, 'wx', 0o600)
+      break
     } catch (error) {
       if (isEEXIST(error)) continue
       throw error
     }
-    let failure: unknown
-    try {
-      await handle.writeFile(bytes)
-      await handle.sync()
-    } catch (error: unknown) {
-      failure = error
-    }
-    try {
-      await handle.close()
-    } catch (error: unknown) {
-      failure = failure === undefined
-        ? error
-        : new AggregateError([failure, error], `failed to write and close migration stage "${path}"`)
-    }
-    if (failure !== undefined) {
-      const writeError = failure instanceof Error
-        ? failure
-        : new Error('migration stage write failed with a non-Error rejection', { cause: failure })
-      try {
-        await internals.fs.rm(path)
-      } catch (cleanupError: unknown) {
-        throw new AggregateError([writeError, cleanupError], `failed to clean migration stage "${path}"`)
+  }
+  const hash = createHash('sha256')
+  let bytes = 0
+  const write = async (chunk: Buffer): Promise<void> => {
+    await handle.writeFile(chunk)
+    hash.update(chunk)
+    bytes += chunk.length
+  }
+  let failure: unknown
+  try {
+    const headerValue = format.encodeHeader(artifact.header, artifact.inheritedEventCount)
+    const header = Buffer.from(`${stringifyJson(headerValue, 'migrated session header')}\n`)
+    await write(compression === 'zstd' ? await compressZstdFrame(header) : header)
+    if (artifact.events.length > 0) {
+      const rows = encodeMigrationRows(artifact, format, signal)
+      if (compression === 'none') {
+        await writeMigrationChunks(rows, write)
+      } else {
+        await new Promise<void>((resolve, reject) => {
+          pipeline(
+            Readable.from(rows, { objectMode: false, highWaterMark: MIGRATION_WORK_CHUNK_BYTES }),
+            createZstdCompress(ZSTD_CHECKSUM_OPTIONS),
+            async (source) => { await writeMigrationChunks(source as AsyncIterable<Buffer>, write) },
+            (error: Error | null | undefined) => {
+              if (error instanceof Error) reject(error)
+              else resolve()
+            },
+          )
+        })
       }
-      throw writeError
     }
-    return path
+    signal?.throwIfAborted()
+    await handle.sync()
+  } catch (error: unknown) {
+    failure = error
   }
+  try {
+    await handle.close()
+  } catch (error: unknown) {
+    failure = failure === undefined
+      ? error
+      : new AggregateError([failure, error], `failed to write and close migration stage "${path}"`)
+  }
+  if (failure !== undefined) {
+    const writeError = failure instanceof Error
+      ? failure
+      : new Error('migration stage write failed with a non-Error rejection', { cause: failure })
+    await removeTemporary(path, writeError, internals)
+    throw writeError
+  }
+  return { path, bytes, digest: hash.digest('hex') }
 }
 
 /** Remove one temporary file without hiding the operation failure that made it disposable. */
@@ -550,31 +843,6 @@ async function removeCommittedTemporary(
   }
 }
 
-async function validatePhysicalCurrent(
-  path: string,
-  compression: JsonlCompression,
-  format: JsonlGenerationFormatAdapter,
-  signal: AbortSignal | undefined,
-  internals: JsonlGenerationInternals,
-): Promise<JsonlPhysicalSnapshot> {
-  const snapshot = await readStableSnapshot(path, signal, internals.fs)
-  const decoded = await decodePhysicalJsonl(snapshot.bytes, compression, signal)
-  if (decoded.torn) throw new Error('staged current session generation has a torn physical tail')
-  const generation = parseGeneration(decoded.bytes)
-  if (storedVersion(generation.header) !== format.currentVersion) {
-    throw new Error(`staged session generation is not current v${format.currentVersion}`)
-  }
-  format.validateCurrent(generation)
-  const headerEnd = decoded.bytes.indexOf(0x0A)
-  /* v8 ignore next -- parseGeneration already required the header newline. */
-  if (headerEnd === -1) throw new Error('empty or header-less session log')
-  return {
-    ...snapshot,
-    headerValue: generation.header,
-    headerRecord: Buffer.from(decoded.bytes.subarray(0, headerEnd + 1)),
-  }
-}
-
 async function publishCurrentExclusive(
   staged: string,
   currentPath: string,
@@ -609,15 +877,12 @@ function asError(error: unknown): Error {
   })
 }
 
-async function reopenExpectedCurrent(
+async function inspectExpectedCurrent<T>(
   currentPath: string,
-  expectedBytes: Buffer,
-  compression: JsonlCompression,
-  format: JsonlGenerationFormatAdapter,
-  signal: AbortSignal | undefined,
   checkCanonicalTargetName: boolean,
   internals: JsonlGenerationInternals,
-): Promise<JsonlPhysicalSnapshot> {
+  inspect: () => Promise<T>,
+): Promise<T> {
   try {
     if (checkCanonicalTargetName) {
       const expectedName = basename(currentPath)
@@ -631,23 +896,16 @@ async function reopenExpectedCurrent(
     }
     const info = await internals.fs.lstat(currentPath)
     if (info.isSymbolicLink() || !info.isFile()) {
-      const kind = info.isSymbolicLink() ? 'symbolic link' : 'non-regular file'
-      throw new Error(`target is a ${kind}`)
+      throw new Error(`target is a ${info.isSymbolicLink() ? 'symbolic link' : 'non-regular file'}`)
     }
-    const snapshot = await validatePhysicalCurrent(currentPath, compression, format, signal, internals)
-    if (snapshot.bytes.length < expectedBytes.length
-      || !snapshot.bytes.subarray(0, expectedBytes.length).equals(expectedBytes)) {
-      throw new Error('target bytes do not begin with the migrated generation')
-    }
-    return snapshot
+    return await inspect()
   } catch (error: unknown) {
-    if (signal?.aborted) signal.throwIfAborted()
     if (isErrnoException(error)) throw error
     throw new JsonlGenerationTargetConflictError(currentPath, asError(error))
   }
 }
 
-function withOverrides(overrides: JsonlGenerationTestOverrides): JsonlGenerationInternals {
+function withOverrides(overrides: JsonlGenerationRuntimeOverrides): JsonlGenerationInternals {
   return {
     ...defaultInternals,
     ...overrides,
@@ -655,6 +913,35 @@ function withOverrides(overrides: JsonlGenerationTestOverrides): JsonlGeneration
   }
 }
 
+async function reopenExpectedCurrent(
+  currentPath: string,
+  staged: StreamedMigrationStage,
+  compression: JsonlCompression,
+  expectedId: string,
+  expectedEventCount: number,
+  verifyCurrentFile: EnsureJsonlGenerationOptions['verifyCurrentFile'],
+  signal: AbortSignal | undefined,
+  checkCanonicalTargetName: boolean,
+  internals: JsonlGenerationInternals,
+): Promise<JsonlPhysicalSnapshot> {
+  return inspectExpectedCurrent(currentPath, checkCanonicalTargetName, internals, async () => {
+    const verified = await verifyCurrentFile(
+      currentPath,
+      compression,
+      expectedId,
+      expectedEventCount,
+      staged,
+      signal,
+    )
+    if (verified.bytes !== staged.bytes || verified.digest !== staged.digest) {
+      throw new Error('target bytes differ from the migrated generation')
+    }
+    const snapshot = await readStableSnapshot(currentPath, signal, internals.fs)
+    const header = readPhysicalHeader(snapshot.bytes, compression, signal)
+    return { ...snapshot, headerValue: header.value, headerRecord: header.record }
+  })
+}
+
 async function ensureCurrent(
   options: EnsureJsonlGenerationOptions,
   internals: JsonlGenerationInternals,
@@ -681,68 +968,87 @@ async function ensureCurrent(
       )
     }
     if (sourceVersion > format.currentVersion) {
-      throw new JsonlGenerationNewerVersionError(sourceVersion, format.currentVersion, storedId(quickHeader.value))
+      throw new JsonlGenerationNewerVersionError(
+        sourceVersion,
+        format.currentVersion,
+        storedId(quickHeader.value),
+      )
     }
     if (sourceVersion === format.currentVersion) {
       return {
         status: 'current',
         version: quickVersion,
         path: sourcePath,
-        snapshot: { ...source, headerValue: quickHeader.value, headerRecord: quickHeader.record },
+        snapshot: {
+          ...source,
+          headerValue: quickHeader.value,
+          headerRecord: quickHeader.record,
+        },
       }
     }
-    const validation = options.validateHistoricalHeader?.(quickHeader.value)
-    if (validation !== undefined) await validation
-
-    const decodedSource = await decodePhysicalJsonl(source.bytes, compression, signal)
-    const parsedSource = parseGeneration(decodedSource.bytes, true)
-    const fromVersion = storedVersion(parsedSource.header)
-    /* v8 ignore next -- both headers come from the same stable physical snapshot. */
-    if (fromVersion !== quickVersion) throw new Error('session format changed within one stable physical snapshot')
-    const sourceFingerprint = fingerprint(source.identity, source.bytes)
-
-    let migrated: JsonlCurrentGeneration
+    let artifact: SessionFormatArtifact
     try {
-      migrated = format.migrate(parsedSource)
+      artifact = await decodeStreamingMigration(
+        source.bytes,
+        compression,
+        format,
+        options.validateHistoricalHeader,
+        signal,
+      )
     } catch (error: unknown) {
       if (format.isUnsupportedMigrationError?.(error) === true) {
-        throw new JsonlGenerationUnsupportedMigrationError(fromVersion, error)
+        throw new JsonlGenerationUnsupportedMigrationError(sourceVersion, error)
       }
       throw error
     }
-    if (storedVersion(migrated.header) !== format.currentVersion) {
-      throw new Error(`format migration returned v${storedVersion(migrated.header)}, expected v${format.currentVersion}`)
+    if (artifact.header.version !== format.currentVersion) {
+      throw new Error(`format migration returned v${artifact.header.version}, expected v${format.currentVersion}`)
     }
-    const logical = encodeLogicalJsonl(migrated)
-    const physical = await encodePhysicalJsonl(logical, migrated, compression)
-    let staged = await writeSyncedTemp(currentPath, suffix, physical, internals)
+
+    await scheduler.yield()
+    signal?.throwIfAborted()
+    const sourceFingerprint = fingerprint(source.identity, source.bytes)
+    const eventCount = artifact.events.length
+    let staged = await writeSyncedTemp(currentPath, suffix, compression, artifact, format, signal, internals)
     let failure: unknown
     try {
-      await validatePhysicalCurrent(staged, compression, format, signal, internals)
+      const verifiedStage = await options.verifyCurrentFile(
+        staged.path,
+        compression,
+        artifact.header.id,
+        eventCount,
+        undefined,
+        signal,
+      )
+      if (verifiedStage.bytes !== staged.bytes || verifiedStage.digest !== staged.digest) {
+        throw new Error('staged session generation changed during verification')
+      }
       await internals.barrier('before-source-check', attempt)
       const beforePublish = await readStableSnapshot(sourcePath, signal, internals.fs)
       if (fingerprint(beforePublish.identity, beforePublish.bytes) !== sourceFingerprint) continue
 
-      const published = await publishCurrentExclusive(staged, currentPath, internals)
-      if (published && internals.platform === 'win32') staged = ''
+      const published = await publishCurrentExclusive(staged.path, currentPath, internals)
+      if (published && internals.platform === 'win32') staged = { ...staged, path: '' }
       await internals.barrier('after-publication', attempt)
       signal?.throwIfAborted()
       const committed = await reopenExpectedCurrent(
         currentPath,
-        physical,
+        staged,
         compression,
-        format,
+        artifact.header.id,
+        eventCount,
+        options.verifyCurrentFile,
         signal,
         !published,
         internals,
       )
-      if (staged !== '') {
-        await removeCommittedTemporary(staged, internals)
-        staged = ''
+      if (staged.path !== '') {
+        await removeCommittedTemporary(staged.path, internals)
+        staged = { ...staged, path: '' }
       }
       return {
         status: 'migrated',
-        fromVersion,
+        fromVersion: sourceVersion,
         toVersion: format.currentVersion,
         path: currentPath,
         sourcePath,
@@ -752,32 +1058,38 @@ async function ensureCurrent(
       failure = error
       throw error
     } finally {
-      if (staged !== '') await removeTemporary(staged, failure, internals)
+      if (staged.path !== '') await removeTemporary(staged.path, failure, internals)
     }
   }
 }
 
 /**
- * Ensure one resolved generation has a current-format successor. Current input reads one
- * coherent physical snapshot, inspects only its independently readable header,
- * invokes no body decoder or migration callback, and returns that snapshot for
- * the immediate body-reading backend hook. Historical input remains unchanged;
- * only a previously absent current filename can be published.
- * @param options - resolved source and target, configured encoding, format adapter, and cancellation.
- * @returns whether the source was already current or which immutable successor was published.
+ * Ensure one resolved generation has a current-format successor before returning.
+ * @param options - resolved source, current target, format adapter, verification, and cancellation.
+ * @returns the current source or the verified and reopened migrated successor.
  */
 export function ensureJsonlGenerationCurrent(
   options: EnsureJsonlGenerationOptions,
 ): Promise<EnsureJsonlGenerationResult> {
-  return ensureCurrent(options, defaultInternals)
+  return defaultGenerationRuntime.ensure(options)
 }
 
-/** Private deterministic filesystem, platform, and race seams for package tests. */
-export const __jsonlGenerationTest = {
-  ensure(
-    options: EnsureJsonlGenerationOptions,
-    overrides: JsonlGenerationTestOverrides,
-  ): Promise<EnsureJsonlGenerationResult> {
-    return ensureCurrent(options, withOverrides(overrides))
-  },
+/**
+ * Create one generation runtime with fixed filesystem and publication dependencies.
+ * @param overrides - deterministic filesystem, platform, and race dependencies.
+ * @returns bound generation operations.
+ */
+export function createJsonlGenerationRuntime(
+  overrides: JsonlGenerationRuntimeOverrides = {},
+): JsonlGenerationRuntime {
+  const internals = withOverrides(overrides)
+  return {
+    readStable: (path, signal) => readStableSnapshot(path, signal, internals.fs),
+    ensure: options => ensureCurrent(options, internals),
+    verify: (path, compression, expectedId, expectedEventCount, expectedPrefix) => verifyCurrentGeneration(
+      path, compression, expectedId, expectedEventCount, internals.fs, expectedPrefix,
+    ),
+  }
 }
+
+const defaultGenerationRuntime = createJsonlGenerationRuntime()

+ 21 - 19
packages/session/session-persistence-jsonl/src/index.ts

@@ -42,6 +42,7 @@ import {
   compressZstdFrame, createZstdFrameDecoder, decompressZstdFrame, decompressZstdPrefix, scanZstdFrames,
 } from './zstd.ts'
 import { ensureDurableDirectoryWin32, publishNewFileWin32 } from './win32.ts'
+import { verifyCurrentGenerationInWorker } from './migration-verifier.ts'
 import {
   ensureJsonlGenerationCurrent,
   JsonlGenerationNewerVersionError,
@@ -180,15 +181,13 @@ class JsonlSessionPersistence extends SessionPersistence {
     this.compression = config.compression ?? DEFAULT_COMPRESSION
     this.generationFormat = {
       currentVersion: sessionFormatCatalog.currentVersion,
-      migrate: (source) => {
-        const decoded = sessionFormatCatalog.decodeRecoverableArtifact(source.header, source.rows)
-        const current = sessionFormatCatalog.migrate(decoded)
-        return sessionFormatCatalog.encodeCurrent(current)
-      },
-      validateCurrent: (candidate) => {
-        const decoded = sessionFormatCatalog.decodeArtifact(candidate.header, candidate.rows)
-        sessionFormatCatalog.migrate(decoded)
-      },
+      createRestore: header => sessionFormatCatalog.createRestore(header, {
+        recovery: 'recoverable',
+        validation: 'transformed',
+      }),
+      encodeHeader: (header, inheritedEventCount) =>
+        sessionFormatCatalog.encodeCurrentHeader(header, inheritedEventCount),
+      encodeEvent: event => sessionFormatCatalog.encodeCurrentEvent(event),
       isUnsupportedMigrationError: (error): error is SessionFormatUnsupportedMigrationError =>
         error instanceof SessionFormatUnsupportedMigrationError,
     }
@@ -397,8 +396,6 @@ class JsonlSessionPersistence extends SessionPersistence {
       }
     }
     const current = await this.ensureCurrentLog(id, signal, selected)
-    /* v8 ignore next -- supplying a resolved generation makes absence unreachable. */
-    if (current === undefined) throw new SessionPersistenceNotFoundError(id)
     return this.decodeStoredLog(
       current.path,
       id,
@@ -411,12 +408,10 @@ class JsonlSessionPersistence extends SessionPersistence {
   /** Select and, when required, publish one immutable current generation. */
   private async ensureCurrentLog(
     id: SessionId,
-    signal?: AbortSignal,
-    resolved?: ResolvedJsonlGeneration,
-  ): Promise<EnsureJsonlGenerationResult | undefined> {
+    signal: AbortSignal | undefined,
+    selected: ResolvedJsonlGeneration,
+  ): Promise<EnsureJsonlGenerationResult> {
     signal?.throwIfAborted()
-    const selected = resolved ?? await this.findLog(id, signal)
-    if (selected === undefined) return undefined
     try {
       return await ensureJsonlGenerationCurrent({
         sourcePath: selected.sourcePath,
@@ -424,6 +419,7 @@ class JsonlSessionPersistence extends SessionPersistence {
         currentPath: selected.currentPath,
         compression: this.compression,
         format: this.generationFormat,
+        verifyCurrentFile: verifyCurrentGenerationInWorker,
         validateHistoricalHeader: headerValue => this.validateSourceIdentity(
           selected,
           headerValue,
@@ -549,8 +545,8 @@ class JsonlSessionPersistence extends SessionPersistence {
     const selected = await this.findLog(id, signal)
     if (selected === undefined) return undefined
     if (selected.sourceVersion === SESSION_FORMAT_VERSION) return selected.sourcePath
-    const current = await this.ensureCurrentLog(id, signal)
-    return current?.path
+    const current = await this.ensureCurrentLog(id, signal, selected)
+    return current.path
   }
 
   /**
@@ -794,8 +790,14 @@ class JsonlSessionPersistence extends SessionPersistence {
       )
     }
     if (result.status === 'unsupported') {
+      const physicalId = String((value as { id?: unknown }).id)
+      let reason = result.reason
+      /* v8 ignore else -- released historical header migrations cannot refuse after physical decoding. */
+      if (result.storedVersion > SESSION_FORMAT_VERSION) {
+        reason = sessionFormatVersionRefusal(physicalId, result.storedVersion)
+      }
       throw new SessionFormatUnsupportedError(
-        `${result.reason} (raw log: ${selected.sourcePath})`,
+        `${reason} (raw log: ${selected.sourcePath})`,
         { kind: 'jsonl', path: selected.sourcePath },
       )
     }

+ 193 - 0
packages/session/session-persistence-jsonl/src/migration-verifier.ts

@@ -0,0 +1,193 @@
+/** Isolated verification for a staged or competing current JSONL generation. */
+
+import { Worker } from 'node:worker_threads'
+import type { WorkerOptions } from 'node:worker_threads'
+import type { JsonlCompression } from './format.ts'
+import type { JsonlExpectedPrefix, JsonlVerifiedGeneration } from './generation.ts'
+
+interface VerificationRequest {
+  readonly path: string
+  readonly compression: JsonlCompression
+  readonly expectedId: string
+  readonly expectedEventCount: number
+  readonly expectedPrefix?: JsonlExpectedPrefix
+}
+
+type VerificationResponse =
+  | { readonly ok: true; readonly result: JsonlVerifiedGeneration }
+  | { readonly ok: false; readonly message: string; readonly stack?: string }
+
+/** Process-wide memory bound for full-generation verification isolates. */
+const MAX_CONCURRENT_VERIFIERS = 2
+
+class VerificationScheduler {
+  private active = 0
+  private readonly waiting: Array<{ grant(): void }> = []
+
+  async run<T>(operation: () => Promise<T>, signal?: AbortSignal): Promise<T> {
+    const permit = this.acquire(signal)
+    if (permit !== undefined) await permit
+    try {
+      signal?.throwIfAborted()
+      return await operation()
+    } finally {
+      this.release()
+    }
+  }
+
+  private acquire(signal?: AbortSignal): Promise<void> | undefined {
+    signal?.throwIfAborted()
+    if (this.active < MAX_CONCURRENT_VERIFIERS) {
+      this.active += 1
+      return
+    }
+    return new Promise<void>((resolve, reject) => {
+      const waiter = {
+        grant: (): void => {
+          signal?.removeEventListener('abort', abort)
+          resolve()
+        },
+      }
+      const abort = (): void => {
+        const index = this.waiting.indexOf(waiter)
+        this.waiting.splice(index, 1)
+        reject(verifierAbortError(signal))
+      }
+      this.waiting.push(waiter)
+      signal?.addEventListener('abort', abort, { once: true })
+    })
+  }
+
+  private release(): void {
+    const next = this.waiting.shift()
+    if (next === undefined) {
+      this.active -= 1
+      return
+    }
+    next.grant()
+  }
+}
+
+const verificationScheduler = new VerificationScheduler()
+
+function workerSpawn(request: VerificationRequest): { readonly entry: string | URL; readonly options: WorkerOptions } {
+  /* v8 ignore next 3 -- built-worker coverage owns the bundled path. */
+  if (!import.meta.url.endsWith('.ts')) {
+    return {
+      entry: new URL('./worker.cjs', import.meta.url),
+      options: { workerData: request, execArgv: [] },
+    }
+  }
+  const workerEntry = new URL('./worker.ts', import.meta.url)
+  const bootstrap = [
+    `import { register as registerEsm } from ${JSON.stringify(import.meta.resolve('tsx/esm/api'))}`,
+    `import { register as registerCjs } from ${JSON.stringify(import.meta.resolve('tsx/cjs/api'))}`,
+    'registerCjs()',
+    'registerEsm()',
+    `await import(${JSON.stringify(workerEntry.href)})`,
+  ].join('\n')
+  return {
+    entry: new URL(`data:text/javascript,${encodeURIComponent(bootstrap)}`),
+    options: {
+      workerData: request,
+      execArgv: [],
+    },
+  }
+}
+
+/**
+ * Verify one current generation in a fresh Worker Thread.
+ * @param path - staged or competing current-generation path.
+ * @param compression - configured physical encoding.
+ * @param expectedId - Session id expected in the decoded header.
+ * @param expectedEventCount - exact logical event count expected after decoding.
+ * @param expectedPrefix - verified physical prefix; an append tail may be present and is not validated.
+ * @param signal - optional cancellation for scheduler wait and Worker execution.
+ * @returns stable physical identity and digest observed by the worker.
+ */
+export function verifyCurrentGenerationInWorker(
+  path: string,
+  compression: JsonlCompression,
+  expectedId: string,
+  expectedEventCount: number,
+  expectedPrefix?: JsonlExpectedPrefix,
+  signal?: AbortSignal,
+): Promise<JsonlVerifiedGeneration> {
+  return verificationScheduler.run(() => runVerificationWorker(
+    path,
+    compression,
+    expectedId,
+    expectedEventCount,
+    expectedPrefix,
+    signal,
+  ), signal)
+}
+
+function runVerificationWorker(
+  path: string,
+  compression: JsonlCompression,
+  expectedId: string,
+  expectedEventCount: number,
+  expectedPrefix?: JsonlExpectedPrefix,
+  signal?: AbortSignal,
+): Promise<JsonlVerifiedGeneration> {
+  signal?.throwIfAborted()
+  const request: VerificationRequest = {
+    path, compression, expectedId, expectedEventCount,
+    ...(expectedPrefix === undefined ? {} : { expectedPrefix }),
+  }
+  const { entry, options } = workerSpawn(request)
+  const worker = new Worker(entry, options)
+  return new Promise((resolve, reject) => {
+    let settled = false
+    const cleanup = (): void => {
+      signal?.removeEventListener('abort', abort)
+    }
+    const fail = (error: Error): void => {
+      /* v8 ignore next -- a late error/exit races only after another terminal callback settled. */
+      if (settled) return
+      settled = true
+      cleanup()
+      void worker.terminate().then(
+        () => { reject(error) },
+        (cleanup: unknown) => {
+          reject(new AggregateError([error, cleanup], 'migration verifier termination failed'))
+        },
+      )
+    }
+    worker.once('message', (value: unknown) => {
+      /* v8 ignore next -- a duplicate message races only after another terminal callback settled. */
+      if (settled) return
+      if (typeof value !== 'object' || value === null || typeof (value as { ok?: unknown }).ok !== 'boolean') {
+        fail(new Error('migration verifier returned an invalid response'))
+        return
+      }
+      const response = value as VerificationResponse
+      if (!response.ok) {
+        const error = new Error(response.message)
+        if (response.stack !== undefined) error.stack = response.stack
+        fail(error)
+        return
+      }
+      settled = true
+      cleanup()
+      void worker.terminate().then(
+        () => { resolve(response.result) },
+        (error: unknown) => { reject(error instanceof Error ? error : new Error(String(error))) },
+      )
+    })
+    worker.once('error', fail)
+    worker.once('exit', (code) => {
+      if (!settled) fail(new Error(`migration verifier exited before reporting a result (code ${code})`))
+    })
+    const abort = (): void => { fail(verifierAbortError(signal)) }
+    signal?.addEventListener('abort', abort, { once: true })
+  })
+}
+
+function verifierAbortError(signal?: AbortSignal): Error {
+  const reason: unknown = signal?.reason
+  return reason instanceof Error
+    ? reason
+    : new Error('migration verifier aborted', { cause: reason })
+}

+ 16 - 0
packages/session/session-persistence-jsonl/src/testing/generation.ts

@@ -0,0 +1,16 @@
+import {
+  createJsonlGenerationRuntime,
+  type JsonlGenerationRuntime,
+  type JsonlGenerationRuntimeOverrides,
+} from '../generation.ts'
+
+/**
+ * Create generation operations with deterministic I/O and race seams for tests.
+ * @param overrides - deterministic filesystem, platform, and race dependencies.
+ * @returns bound generation operations.
+ */
+export function createJsonlGenerationTestRuntime(
+  overrides: JsonlGenerationRuntimeOverrides = {},
+): JsonlGenerationRuntime {
+  return createJsonlGenerationRuntime(overrides)
+}

+ 56 - 0
packages/session/session-persistence-jsonl/src/worker.ts

@@ -0,0 +1,56 @@
+/** Worker entry for current-generation physical and logical verification. */
+
+import { parentPort, workerData } from 'node:worker_threads'
+import { verifyJsonlCurrentGeneration } from './generation.ts'
+import type { JsonlExpectedPrefix } from './generation.ts'
+import type { JsonlCompression } from './format.ts'
+
+interface VerificationRequest {
+  readonly path: string
+  readonly compression: JsonlCompression
+  readonly expectedId: string
+  readonly expectedEventCount: number
+  readonly expectedPrefix?: JsonlExpectedPrefix
+}
+
+function parseRequest(value: unknown): VerificationRequest {
+  if (typeof value !== 'object' || value === null) throw new Error('migration verifier request must be an object')
+  const request = value as Partial<VerificationRequest>
+  if (typeof request.path !== 'string'
+    || request.compression !== 'none' && request.compression !== 'zstd'
+    || typeof request.expectedId !== 'string'
+    || !Number.isSafeInteger(request.expectedEventCount)
+    || (request.expectedEventCount as number) < 0
+    || request.expectedPrefix !== undefined
+      && (!Number.isSafeInteger(request.expectedPrefix.bytes)
+        || request.expectedPrefix.bytes < 0
+        || !/^[0-9a-f]{64}$/.test(request.expectedPrefix.digest))) {
+    throw new Error('migration verifier request is malformed')
+  }
+  return request as VerificationRequest
+}
+
+if (parentPort === null) throw new Error('migration verifier requires a parent port')
+const port = parentPort
+
+const request = parseRequest(workerData)
+
+async function verify(): Promise<void> {
+  try {
+    const result = await verifyJsonlCurrentGeneration(
+      request.path,
+      request.compression,
+      request.expectedId,
+      request.expectedEventCount,
+      request.expectedPrefix,
+    )
+    port.postMessage({ ok: true, result })
+  } catch (error: unknown) {
+    const failure = error instanceof Error ? error : new Error(String(error))
+    port.postMessage({ ok: false, message: failure.message, stack: failure.stack })
+  } finally {
+    port.close()
+  }
+}
+
+void verify()

+ 49 - 0
packages/session/session-persistence-jsonl/tests/built-migration-worker.e2e.ts

@@ -0,0 +1,49 @@
+import { existsSync } from 'node:fs'
+import { join } from 'node:path'
+import { fileURLToPath } from 'node:url'
+import { execa } from 'execa'
+import { describe, expect, it } from 'vitest'
+
+const packageRoot = fileURLToPath(new URL('..', import.meta.url))
+const built = ['lib/index.js', 'lib/worker.cjs']
+  .every(path => existsSync(join(packageRoot, path)))
+
+describe.skipIf(!built)('built migration verifier (plain node)', () => {
+  it('publishes a historical generation through the bundled worker', async () => {
+    const script = `
+      import { mkdir, mkdtemp, readFile, rm, writeFile } from 'node:fs/promises'
+      import { tmpdir } from 'node:os'
+      import { join } from 'node:path'
+      import { Context } from '@deepseek-ai/cordis'
+      import JsonlSessionPersistence from '@deepseek-ai/dsh-session-persistence-jsonl'
+
+      const root = await mkdtemp(join(tmpdir(), 'dsh-built-migration-'))
+      const id = 'built-migration-worker'
+      const directory = join(root, '_no-cwd', id)
+      const ctx = new Context()
+      try {
+        await mkdir(directory, { recursive: true })
+        await writeFile(join(directory, 'session.jsonl'), JSON.stringify({
+          type: 'session', version: 0, id, createdAt: 1, delegationDepth: 0,
+        }) + '\\n')
+        await ctx.plugin(JsonlSessionPersistence, { root, compression: 'none' })
+        const handle = await ctx.sessionPersistence.open(id, 'read')
+        await handle.close()
+        await ctx.sessionPersistence.flush()
+        const header = JSON.parse((await readFile(join(directory, 'session.v2.jsonl'), 'utf8')).trim())
+        console.log(JSON.stringify({ id: header.id, version: header.version }))
+      } finally {
+        await ctx.fiber.dispose()
+        await rm(root, { recursive: true, force: true })
+      }
+    `
+    const { exitCode, stdout, stderr } = await execa(
+      process.execPath,
+      ['--input-type=module', '-e', script],
+      { cwd: packageRoot, stdin: 'ignore', timeout: 30_000, killSignal: 'SIGKILL', reject: false },
+    )
+
+    expect(exitCode, `stderr:\n${stderr}`).toBe(0)
+    expect(JSON.parse(stdout.trim())).toEqual({ id: 'built-migration-worker', version: 2 })
+  })
+})

Fișier diff suprimat deoarece este prea mare
+ 552 - 134
packages/session/session-persistence-jsonl/tests/generation.spec.ts


+ 43 - 36
packages/session/session-persistence-jsonl/tests/jsonl.spec.ts

@@ -19,7 +19,6 @@ import {
 import { runLiveWritePathContract } from '../../session-persistence/tests/live-write-contract.ts'
 import { LIVE_WRITE_BATCH_MAX_DELAY_MS, type JsonlSessionHandle } from '../src/storage.ts'
 import SessionStore from '@deepseek-ai/dsh-session'
-import { releasedV1SessionFormatCodec } from '@deepseek-ai/dsh-session-format-v0-to-v1'
 
 const statRace = vi.hoisted(() => ({
   path: undefined as string | undefined,
@@ -154,15 +153,15 @@ function releasedV1PackedPhysicalLog(header: SessionHeader): string {
         : {}),
     })),
   ]
-  const encoded = releasedV1SessionFormatCodec.encodeArtifact({
-    header: { ...header, version: 1, delegationDepth: header.delegationDepth ?? 0 },
-    inheritedEventCount: 0,
-    events,
-  } as never, { packChunks: true })
-  if (!encoded.rows.some(row => row['type'] === 'text-chunks')) {
-    throw new Error('released v1 test fixture did not produce a packed text row')
+  const packed = {
+    type: 'text-chunks',
+    seq0: 4,
+    time0: 3,
+    data: { turn: 1, step: 1, index: 0, dt: [0, 0], texts: ['hello', '', ''] },
   }
-  return [encoded.header, ...encoded.rows].map(row => JSON.stringify(row)).join('\n') + '\n'
+  const rows = [...events.slice(0, 4), packed, ...events.slice(7)]
+  return [{ ...releasedV0Header(header), version: 1 }, ...rows]
+    .map(row => JSON.stringify(row)).join('\n') + '\n'
 }
 
 /** Create + append + close: persist one whole log through the write handle. */
@@ -273,8 +272,8 @@ describe('JsonlSessionPersistence: format helpers', () => {
   it('ignores retired-header checks for non-object values', () => {
     expect(() => { assertNoRetiredHeaderFields(null) }).not.toThrow()
     expect(() => { assertNoRetiredHeaderFields('header') }).not.toThrow()
-    expect(() => scanLog(Buffer.from('42\n'))).toThrow(/session header/)
-    expect(() => scanLog(Buffer.from('null\n'))).toThrow(/session header/)
+    expect(() => scanLog(Buffer.from('42\n'))).toThrow(/first line is not a JSON object/)
+    expect(() => scanLog(Buffer.from('null\n'))).toThrow(/first line is not a JSON object/)
     expect(() => scanLog(Buffer.from(
       `${JSON.stringify({ type: 'session', version: 2, id: 123 })}\n`,
     ))).toThrow(/first line is not a session header/)
@@ -522,7 +521,7 @@ describe('JsonlSessionPersistence: stored-format refusals', () => {
     // Simulate the race: the header-only stat sees nothing although the full
     // log is present and readable.
     vi.spyOn(ctx.sessionPersistence, 'stat').mockResolvedValue(undefined)
-    const handle = await ctx.sessionPersistence.open(m.id, 'read')
+    const handle = await ctx.sessionPersistence.open(m.id, 'read', { signal: new AbortController().signal })
     try {
       expect(handle.header).toMatchObject({ id: m.id, cwd: '/work' })
       expect(await handle.read()).toEqual(oneTurnLog())
@@ -612,7 +611,6 @@ describe('JsonlSessionPersistence: immutable format generations', () => {
       meta: { ...header, delegationDepth: 0 },
       events: oneTurnLog(),
     })
-
     expect(await readFile(sourcePath)).toEqual(source)
     const current = (await readFile(currentPath, 'utf8')).trimEnd().split('\n')
     expect(JSON.parse(current[0] as string)).toMatchObject({
@@ -720,14 +718,6 @@ describe('JsonlSessionPersistence: immutable format generations', () => {
     expect(await readFile(sourcePath, 'utf8')).toBe(`${JSON.stringify(releasedV0Header(header))}\n`)
   })
 
-  it('reports no current generation when ensure-current finds no stored id', async () => {
-    const storage = ctx.sessionPersistence as unknown as {
-      ensureCurrentLog(id: SessionId, signal?: AbortSignal): Promise<unknown>
-    }
-
-    await expect(storage.ensureCurrentLog(SessionId('missing-generation'))).resolves.toBeUndefined()
-  })
-
   it('opens the migrated successor for append while retaining the historical source', async () => {
     const header = meta('released-v0-write', '/work')
     const sourcePath = historicalLogPath(root, header.cwd, header.id)
@@ -771,6 +761,7 @@ describe('JsonlSessionPersistence: immutable format generations', () => {
     const handle = await ctx.sessionPersistence.create(header, {
       inheritedEventCount: SessionLogOffset(0),
     })
+    expect(handle.inheritedEventCount).toBe(SessionLogOffset(0))
     await handle.append([{
       type: 'session/end-seed', seq: SessionSeq(0), time: 1, data: { inherited: true },
     }])
@@ -1002,16 +993,16 @@ describe('JsonlSessionPersistence: durability and crash semantics', () => {
     const m = meta('opening-claim', '/work')
     await writeLog(ctx.sessionPersistence, m, oneTurnLog())
     const service = ctx.sessionPersistence as unknown as {
-      ensureCurrentLog: (
+      requireStoredLog: (
         id: SessionId,
         signal?: AbortSignal,
         resolved?: unknown,
       ) => Promise<unknown>
     }
-    const original = service.ensureCurrentLog.bind(service)
+    const original = service.requireStoredLog.bind(service)
     const gate = Promise.withResolvers<undefined>()
     const entered = Promise.withResolvers<undefined>()
-    vi.spyOn(service, 'ensureCurrentLog').mockImplementationOnce(async (...args) => {
+    vi.spyOn(service, 'requireStoredLog').mockImplementationOnce(async (...args) => {
       entered.resolve(undefined)
       await gate.promise
       return original(...args)
@@ -1637,18 +1628,34 @@ describe('JsonlSessionPersistence: scanLog unit', () => {
     expect(() => new SessionLogScanner(Buffer.from(`${header}\n${header}\n`))).toThrow(/header-less/)
   })
 
+  it('fails immediately on an invalid row in strict scanner mode', () => {
+    const header = Buffer.from(`${JSON.stringify(toHeaderLine(meta('scanner-strict')))}\n`)
+    const scanner = new SessionLogScanner(header, 'strict')
+    expect(() => { scanner.write(Buffer.from('null\n')) }).toThrow(/invalid committed event/)
+  })
+
+  it('expands valid stored provenance ranges', () => {
+    const log = [
+      JSON.stringify(toHeaderLine(meta('scanner-provenance'))),
+      JSON.stringify({ type: 'external', seq: 0, time: 2, data: null, sourceEventSeqs: [] }),
+      '',
+    ].join('\n')
+    const restored = scanLog(Buffer.from(log)).events[0]
+    expect(restored !== undefined && 'sourceEventSeqs' in restored ? restored.sourceEventSeqs : undefined).toEqual([])
+  })
+
   it('requires the tagged inherited cut to agree with the v2 header lineage', () => {
     const seeded = { ...meta('scanner-seeded-cut'), isSeeded: true }
     const seededHeader = JSON.stringify(toHeaderLine(seeded, SessionLogOffset(0)))
     expect(() => scanLog(Buffer.from(`${seededHeader}\n`)))
-      .toThrow(/seeded v2 header lacks an inherited end-seed marker/)
+      .toThrow(/seeded Session lacks an inherited end-seed marker/)
 
     const unseededHeader = JSON.stringify(toHeaderLine(meta('scanner-unseeded-cut')))
     const inheritedMarker = JSON.stringify({
       type: 'session/end-seed', seq: 0, time: 1, data: { inherited: true },
     })
     expect(() => scanLog(Buffer.from(`${unseededHeader}\n${inheritedMarker}\n`)))
-      .toThrow(/unseeded v2 header contains an inherited end-seed marker/)
+      .toThrow(/unseeded Session contains an inherited end-seed marker/)
   })
 
   it('handles empty writes, boundary newlines, torn fragments, and scanner completion', () => {
@@ -1681,7 +1688,7 @@ describe('JsonlSessionPersistence: scanLog unit', () => {
     expect(() => { committed.write(Buffer.from([
       JSON.stringify({ type: 'turn/end', seq: SessionSeq(1), time: 2, data: { turn: 1, reason: { kind: 'completed' } } }),
       '',
-    ].join('\n'))) }).toThrow(/seq gap in committed region/)
+    ].join('\n'))) }).toThrow(/seq gap/)
   })
 
   it('incrementally scans records split across reusable decoder chunks', () => {
@@ -1806,22 +1813,22 @@ describe('JsonlSessionPersistence: scanLog unit', () => {
     ].join('\n') + '\n'
     // A turn/end exists, so the prefix up to it is committed — but it has a hole.
     // Truncating it would silently drop committed data → unloadable.
-    expect(() => scanLog(Buffer.from(log))).toThrow(/seq gap in committed region/)
+    expect(() => scanLog(Buffer.from(log))).toThrow(/seq gap/)
   })
 
   it('rejects malformed records before a later committed turn/end', () => {
     const corruptRecords = [
-      '{not json',
-      'null',
-      JSON.stringify({ type: 'assistant/message', sourceEventSeqs: [0], data: {} }),
-    ]
-    for (const record of corruptRecords) {
+      ['{not json', /unparsable committed event/],
+      ['null', /invalid committed event/],
+      [JSON.stringify({ type: 'assistant/message', sourceEventSeqs: [0], data: {} }), /invalid committed event/],
+    ] as const
+    for (const [record, message] of corruptRecords) {
       const log = [
         JSON.stringify({ type: 'session', version: SESSION_FORMAT_VERSION, id: 'c', createdAt: 1, isSeeded: false, delegationDepth: 0 }),
         record,
         JSON.stringify({ type: 'turn/end', seq: SessionSeq(1), time: 2, data: { turn: 1, reason: { kind: 'completed' } } }),
       ].join('\n') + '\n'
-      expect(() => scanLog(Buffer.from(log))).toThrow(/unparsable committed event/)
+      expect(() => scanLog(Buffer.from(log))).toThrow(message)
     }
   })
 
@@ -1953,7 +1960,7 @@ describe('JsonlSessionPersistence: nested v2 Assistant streams', () => {
       JSON.stringify({ type: 'text-chunks', seq0: 1, time0: 2, data: { turn: 1, step: 1, index: 0, dt: [1, 1], texts: ['a', 'b', 'c'] } }),
       JSON.stringify({ type: 'turn/end', seq: SessionSeq(4), time: 5, data: { turn: 1, reason: { kind: 'completed' } } }),
     ].join('\n') + '\n'
-    expect(() => scanLog(Buffer.from(logText))).toThrow(/seq gap in committed region/)
+    expect(() => scanLog(Buffer.from(logText))).toThrow(/lacks .*seq/)
   })
 
   it('scanLog treats a malformed removed packed row as a committed seq hole', () => {
@@ -1966,7 +1973,7 @@ describe('JsonlSessionPersistence: nested v2 Assistant streams', () => {
       JSON.stringify({ type: 'text-chunks', seq0: 0, time0: 1, data: { turn: 1, step: 1, index: 0, dt: [], texts: ['a', 'b'] } }),
       JSON.stringify({ type: 'turn/end', seq: SessionSeq(2), time: 3, data: { turn: 1, reason: { kind: 'completed' } } }),
     ].join('\n') + '\n'
-    expect(() => scanLog(Buffer.from(logText))).toThrow(/seq gap in committed region/)
+    expect(() => scanLog(Buffer.from(logText))).toThrow(/lacks .*seq/)
   })
 
   it('scanLog: a packed row with a mid-run seq gap after the last turn/end drops the whole row', () => {

+ 230 - 0
packages/session/session-persistence-jsonl/tests/migration-verifier.spec.ts

@@ -0,0 +1,230 @@
+import { afterEach, describe, expect, it, vi } from 'vitest'
+import { verifyCurrentGenerationInWorker } from '../src/migration-verifier.ts'
+
+const state = vi.hoisted(() => ({ workers: [] as unknown[] }))
+
+vi.mock('node:worker_threads', () => ({
+  Worker: class {
+    readonly listeners = new Map<string, (value: never) => void>()
+    readonly terminate = vi.fn<() => Promise<number>>(() => Promise.resolve(0))
+
+    constructor(readonly entry: string | URL, readonly options: unknown) {
+      state.workers.push(this)
+    }
+
+    once(event: string, listener: (value: never) => void): this {
+      this.listeners.set(event, listener)
+      return this
+    }
+
+    emit(event: string, value: unknown): void {
+      this.listeners.get(event)?.(value as never)
+    }
+  },
+}))
+
+interface FakeWorker {
+  readonly entry: string | URL
+  readonly options: { readonly workerData: unknown }
+  readonly terminate: ReturnType<typeof vi.fn<() => Promise<number>>>
+  emit(event: string, value: unknown): void
+}
+
+function worker(index = 0): FakeWorker {
+  const candidate = state.workers[index]
+  if (candidate === undefined) throw new Error('verification did not create a Worker')
+  return candidate as FakeWorker
+}
+
+const result = {
+  identity: { dev: 1n, ino: 2n, size: 3n, mtimeNs: 4n, ctimeNs: 5n },
+  bytes: 3,
+  digest: 'digest',
+}
+
+afterEach(() => {
+  state.workers.length = 0
+})
+
+describe('migration verifier Worker lifecycle', () => {
+  it('resolves only after terminating a successful Worker', async () => {
+    const verification = verifyCurrentGenerationInWorker('/stage', 'none', 'session', 2)
+    const instance = worker()
+    expect(instance.options.workerData).toEqual({
+      path: '/stage', compression: 'none', expectedId: 'session', expectedEventCount: 2,
+    })
+    instance.emit('message', { ok: true, result })
+
+    await expect(verification).resolves.toEqual(result)
+    expect(instance.terminate).toHaveBeenCalledOnce()
+  })
+
+  it('reconstructs a Worker-reported error', async () => {
+    const verification = verifyCurrentGenerationInWorker('/stage', 'zstd', 'session', 0)
+    worker().emit('message', { ok: false, message: 'invalid stage', stack: 'worker stack' })
+
+    await expect(verification).rejects.toMatchObject({ message: 'invalid stage', stack: 'worker stack' })
+  })
+
+  it('accepts an error response without a stack', async () => {
+    const verification = verifyCurrentGenerationInWorker('/stage', 'none', 'session', 0)
+    worker().emit('message', { ok: false, message: 'invalid stage' })
+    await expect(verification).rejects.toThrow('invalid stage')
+  })
+
+  it.each([
+    ['invalid response', 'message', null, /invalid response/],
+    ['non-object response', 'message', 'invalid', /invalid response/],
+    ['missing discriminator', 'message', {}, /invalid response/],
+    ['invalid discriminator', 'message', { ok: 'yes' }, /invalid response/],
+    ['worker error', 'error', new Error('worker failed'), /worker failed/],
+    ['early exit', 'exit', 7, /code 7/],
+  ])('rejects an %s', async (_name, event, value, expected) => {
+    const verification = verifyCurrentGenerationInWorker('/stage', 'none', 'session', 0)
+    worker().emit(event, value)
+    await expect(verification).rejects.toThrow(expected)
+  })
+
+  it('aggregates termination failure after a Worker failure', async () => {
+    const verification = verifyCurrentGenerationInWorker('/stage', 'none', 'session', 0)
+    const instance = worker()
+    instance.terminate.mockRejectedValueOnce(new Error('terminate failed'))
+    instance.emit('error', new Error('worker failed'))
+
+    await expect(verification).rejects.toBeInstanceOf(AggregateError)
+  })
+
+  it('rejects a successful result when termination fails', async () => {
+    const verification = verifyCurrentGenerationInWorker('/stage', 'none', 'session', 0)
+    const instance = worker()
+    instance.terminate.mockRejectedValueOnce('terminate failed')
+    instance.emit('message', { ok: true, result })
+
+    await expect(verification).rejects.toThrow('terminate failed')
+  })
+
+  it('preserves an Error from successful-result termination', async () => {
+    const verification = verifyCurrentGenerationInWorker('/stage', 'none', 'session', 0)
+    const instance = worker()
+    instance.terminate.mockRejectedValueOnce(new Error('terminate failed'))
+    instance.emit('message', { ok: true, result })
+
+    await expect(verification).rejects.toThrow('terminate failed')
+  })
+
+  it('ignores terminal signals after a result settles', async () => {
+    const verification = verifyCurrentGenerationInWorker('/stage', 'none', 'session', 0)
+    const instance = worker()
+    instance.emit('message', { ok: true, result })
+    instance.emit('error', new Error('late error'))
+    instance.emit('exit', 1)
+    instance.emit('message', null)
+
+    await expect(verification).resolves.toEqual(result)
+    expect(instance.terminate).toHaveBeenCalledOnce()
+  })
+
+  it('starts at most two verification Workers concurrently', async () => {
+    const first = verifyCurrentGenerationInWorker('/first', 'none', 'session', 0)
+    const second = verifyCurrentGenerationInWorker('/second', 'none', 'session', 0)
+    const third = verifyCurrentGenerationInWorker('/third', 'none', 'session', 0)
+    expect(state.workers).toHaveLength(2)
+
+    worker(0).emit('message', { ok: true, result })
+    await first
+    await vi.waitFor(() => { expect(state.workers).toHaveLength(3) })
+
+    worker(1).emit('message', { ok: true, result })
+    worker(2).emit('message', { ok: true, result })
+    await expect(Promise.all([second, third])).resolves.toEqual([result, result])
+  })
+
+  it('hands a released permit directly to the oldest waiter', async () => {
+    const first = verifyCurrentGenerationInWorker('/first', 'none', 'session', 0)
+    const second = verifyCurrentGenerationInWorker('/second', 'none', 'session', 0)
+    const third = verifyCurrentGenerationInWorker('/third', 'none', 'session', 0)
+    let fourth: Promise<typeof result> | undefined
+    worker(0).terminate.mockReturnValueOnce({
+      then(onFulfilled: (value: number) => unknown) {
+        onFulfilled(0)
+        queueMicrotask(() => {
+          fourth = verifyCurrentGenerationInWorker('/fourth', 'none', 'session', 0)
+        })
+        return Promise.resolve()
+      },
+    } as unknown as Promise<number>)
+
+    worker(0).emit('message', { ok: true, result })
+    await first
+    await vi.waitFor(() => { expect(state.workers).toHaveLength(3) })
+    expect(worker(2).options.workerData).toMatchObject({ path: '/third' })
+
+    worker(1).emit('message', { ok: true, result })
+    await second
+    await vi.waitFor(() => { expect(state.workers).toHaveLength(4) })
+    expect(worker(3).options.workerData).toMatchObject({ path: '/fourth' })
+    if (fourth === undefined) throw new Error('fourth verification was not scheduled')
+
+    worker(2).emit('message', { ok: true, result })
+    worker(3).emit('message', { ok: true, result })
+    await expect(Promise.all([third, fourth])).resolves.toEqual([result, result])
+  })
+
+  it('removes an aborted waiter without starting another Worker', async () => {
+    const first = verifyCurrentGenerationInWorker('/first', 'none', 'session', 0)
+    const second = verifyCurrentGenerationInWorker('/second', 'none', 'session', 0)
+    const controller = new AbortController()
+    const reason = new Error('queued verification cancelled')
+    const queued = verifyCurrentGenerationInWorker(
+      '/queued', 'none', 'session', 0, undefined, controller.signal,
+    )
+
+    controller.abort(reason)
+    await expect(queued).rejects.toBe(reason)
+    expect(state.workers).toHaveLength(2)
+
+    worker(0).emit('message', { ok: true, result })
+    worker(1).emit('message', { ok: true, result })
+    await expect(Promise.all([first, second])).resolves.toEqual([result, result])
+    expect(state.workers).toHaveLength(2)
+  })
+
+  it('terminates an active Worker before rejecting cancellation', async () => {
+    const controller = new AbortController()
+    const reason = new Error('active verification cancelled')
+    const verification = verifyCurrentGenerationInWorker(
+      '/stage', 'none', 'session', 0, undefined, controller.signal,
+    )
+    const instance = worker()
+    let finishTermination: ((value: number) => void) | undefined
+    instance.terminate.mockReturnValueOnce(new Promise((resolve) => {
+      finishTermination = resolve
+    }))
+    let settled = false
+    void verification.then(
+      () => { settled = true },
+      () => { settled = true },
+    )
+
+    controller.abort(reason)
+    expect(instance.terminate).toHaveBeenCalledOnce()
+    await Promise.resolve()
+    expect(settled).toBe(false)
+
+    finishTermination?.(0)
+    await expect(verification).rejects.toBe(reason)
+  })
+
+  it('wraps a non-Error active cancellation reason', async () => {
+    const controller = new AbortController()
+    const verification = verifyCurrentGenerationInWorker(
+      '/stage', 'none', 'session', 0, undefined, controller.signal,
+    )
+
+    controller.abort('cancelled')
+    await expect(verification).rejects.toMatchObject({
+      message: 'migration verifier aborted',
+      cause: 'cancelled',
+    })
+  })
+})

+ 0 - 1
packages/session/session-persistence-jsonl/tests/zstd.spec.ts

@@ -428,7 +428,6 @@ describe('JsonlSessionPersistence: default Zstandard encoding', () => {
       meta: { ...header, delegationDepth: 0 },
       events: oneTurnLog(),
     })
-
     expect(await readFile(sourcePath)).toEqual(source)
     const current = (await decodeCompleteFrames(await readFile(currentPath))).toString().split('\n')
     expect(JSON.parse(current[0] as string)).toMatchObject({

+ 3 - 0
packages/session/session-persistence-jsonl/tsconfig.json

@@ -17,6 +17,9 @@
     {
       "path": "../../../vendor/schemastery"
     },
+    {
+      "path": "../../llm/llm"
+    },
     {
       "path": "../../core/session"
     },

+ 25 - 0
packages/session/session-persistence-jsonl/tsdown.config.ts

@@ -0,0 +1,25 @@
+import { defineConfig } from 'tsdown'
+
+/** Build the backend and its path-loaded verifier as separate bundles. */
+export default defineConfig([
+  {
+    entry: ['lib/types/index.js'],
+    outDir: 'lib',
+    format: ['esm'],
+    platform: 'node',
+    target: 'es2024',
+    fixedExtension: false,
+    dts: false,
+    clean: false,
+  },
+  {
+    entry: ['lib/types/worker.js'],
+    outDir: 'lib',
+    format: ['cjs'],
+    platform: 'node',
+    target: 'es2024',
+    fixedExtension: false,
+    dts: false,
+    clean: false,
+  },
+])

+ 3 - 0
pnpm-lock.yaml

@@ -7370,6 +7370,9 @@ importers:
 
   packages/session/session-persistence-jsonl:
     dependencies:
+      '@deepseek-ai/dsh-llm':
+        specifier: workspace:^
+        version: link:../../llm/llm
       '@deepseek-ai/dsh-session-format':
         specifier: workspace:^
         version: link:../session-format

+ 3 - 0
scripts/check-workspace-constraints.ts

@@ -156,6 +156,9 @@ const packageFileExtras: Readonly<Record<string, readonly string[]>> = {
   // The Web Host mounts the default-off settings owner independently of each
   // Agent-scoped delegation-tool instance.
   '@deepseek-ai/dsh-tool-subagent': ['lib/model-selection-settings.js'],
+  // The JSONL backend resolves its private verification Worker relative to
+  // import.meta.url; it is shipped without a public package subpath.
+  '@deepseek-ai/dsh-session-persistence-jsonl': ['lib/worker.cjs'],
   // The argv-prefix runner entry ships beside the lib as its own bundle;
   // sandbox-local resolves it through the package's ./runner export. tsdown
   // also shares its generated FFI code through a hashed runtime chunk.

+ 1 - 0
scripts/run-gates.ts

@@ -791,6 +791,7 @@ function builtBinSmokeGate(needs: string[] = ['build']): Gate {
     // unbuilt, so these files self-skip there.
     'packages/workflow/workflow-worker-thread/tests/built-worker.e2e.ts',
     'packages/code-runtime/code-runtime-worker-thread/tests/built-lib.e2e.ts',
+    'packages/session/session-persistence-jsonl/tests/built-migration-worker.e2e.ts',
     'packages/lsp/lsp-stdio/tests/built-lib.e2e.ts',
   ], {
     label: 'built-bin smoke',

Unele fișiere nu au fost afișate deoarece prea multe fișiere au fost modificate în acest diff