From 353226246e9c866b02d580e404990973e28dd552 Mon Sep 17 00:00:00 2001 From: imccyu Date: Wed, 5 Aug 2026 17:27:02 +0800 Subject: [PATCH] perf(jsonl): add reusable zstd frame decoders --- .../src/zstd-private-decoder.ts | 153 ++++++++++++++++++ .../src/zstd-public-decoder.ts | 38 +++++ .../session-persistence-jsonl/src/zstd.ts | 40 ++++- 3 files changed, 230 insertions(+), 1 deletion(-) create mode 100644 packages/session-persistence/session-persistence-jsonl/src/zstd-private-decoder.ts create mode 100644 packages/session-persistence/session-persistence-jsonl/src/zstd-public-decoder.ts diff --git a/packages/session-persistence/session-persistence-jsonl/src/zstd-private-decoder.ts b/packages/session-persistence/session-persistence-jsonl/src/zstd-private-decoder.ts new file mode 100644 index 0000000000..6f2ac095b3 --- /dev/null +++ b/packages/session-persistence/session-persistence-jsonl/src/zstd-private-decoder.ts @@ -0,0 +1,153 @@ +/** + * Node-private synchronous Zstandard frame decoder optimization. + * @module dsh-session-persistence-jsonl/zstd-private-decoder + */ + +import { constants as bufferConstants } from 'node:buffer' +import { createZstdDecompress } from 'node:zlib' +import type { ZstdFrameDecoder, ZstdFrameRange } from './zstd.ts' + +const DECODE_CHUNK_SIZE = 1024 * 1024 + +interface NodeZstdPrivateHandle { + writeSync( + flushFlag: number, + input: Buffer, + inputOffset: number, + inputLength: number, + output: Buffer, + outputOffset: number, + outputLength: number, + ): void +} + +type NodeZstdPrivateWriteState = Uint32Array & { 0: number; 1: number } + +interface NodeZstdPrivateState { + _handle: NodeZstdPrivateHandle | null + _writeState: NodeZstdPrivateWriteState + _defaultFlushFlag: number +} + +type NodeZstdPrivateStream = ReturnType & NodeZstdPrivateState + +/** Return the stream with its observed private Node contract, or reject that optimization. */ +function privateZstdStream(stream: ReturnType): NodeZstdPrivateStream | undefined { + const candidate = stream as unknown as Partial + const handle = candidate._handle + if ( + typeof handle !== 'object' || handle === null + || typeof (handle as { writeSync?: unknown }).writeSync !== 'function' + || !(candidate._writeState instanceof Uint32Array) + || candidate._writeState.length < 2 + || typeof candidate._defaultFlushFlag !== 'number' + ) return undefined + return stream as NodeZstdPrivateStream +} + +/** + * Synchronous multi-frame decoder backed by one Node Zstd stream handle. Node + * exposes synchronous decoding only as a one-shot API, so this adapter uses + * the stream's private handle contract to reuse its native context and output + * chunks across frames. + */ +export class NodePrivateZstdFrameDecoder implements ZstdFrameDecoder { + private readonly output = Buffer.allocUnsafe(DECODE_CHUNK_SIZE) + private decoderError?: Error + private started = false + private closed = false + + private constructor(private readonly stream: NodeZstdPrivateStream) { + this.stream.on('error', (error: Error) => { + this.decoderError ??= error + }) + } + + /** + * Create the optimized decoder when this Node release exposes the expected + * private stream shape. + * @returns a shared decoder, or `undefined` when callers must use the public fallback. + */ + static create(): NodePrivateZstdFrameDecoder | undefined { + const stream = createZstdDecompress({ chunkSize: DECODE_CHUNK_SIZE }) + const privateStream = privateZstdStream(stream) + if (privateStream !== undefined) return new NodePrivateZstdFrameDecoder(privateStream) + stream.close() + return undefined + } + + /** @inheritdoc */ + public *decode(source: Buffer, frames: readonly ZstdFrameRange[]): Generator { + if (this.started) throw new Error('Zstandard frame decoder was already started') + if (this.closed) throw new Error('cannot start a closed Zstandard frame decoder') + this.started = true + try { + for (const frame of frames) { + try { + yield this.decodeFrame(source.subarray(frame.start, frame.end)) + } catch (error) { + throw new Error(`corrupt Zstandard session log: frame at byte ${frame.start} failed validation`, { + cause: error, + }) + } + } + } finally { + this.close() + } + } + + /** Decode one frame; its returned scratch view remains valid until the next call. */ + private decodeFrame(input: Buffer): Buffer { + const handle = this.stream._handle + if (this.closed || handle === null) throw new Error('cannot decode with a closed Zstandard frame decoder') + + let inputOffset = 0 + let inputRemaining = input.length + let outputBytes = 0 + const fullChunks: Buffer[] = [] + for (;;) { + handle.writeSync( + this.stream._defaultFlushFlag, + input, + inputOffset, + inputRemaining, + this.output, + 0, + DECODE_CHUNK_SIZE, + ) + if (this.decoderError !== undefined) throw this.decoderError + + const outputAfter = this.stream._writeState[0] + const inputAfter = this.stream._writeState[1] + const consumed = inputRemaining - inputAfter + const produced = DECODE_CHUNK_SIZE - outputAfter + if (produced > 0) { + outputBytes += produced + if (outputBytes > bufferConstants.MAX_LENGTH) { + throw new Error(`Zstandard frame output exceeds ${bufferConstants.MAX_LENGTH} bytes`) + } + } + + if (outputAfter !== 0) { + if (inputAfter !== 0) throw new Error('Zstandard frame decoder left trailing input') + const finalChunk = this.output.subarray(0, produced) + if (fullChunks.length === 0) return finalChunk + if (produced > 0) fullChunks.push(Buffer.from(finalChunk)) + const [onlyChunk] = fullChunks + return fullChunks.length === 1 && onlyChunk !== undefined + ? onlyChunk + : Buffer.concat(fullChunks, outputBytes) + } + fullChunks.push(Buffer.from(this.output)) + inputOffset += consumed + inputRemaining = inputAfter + } + } + + /** @inheritdoc */ + close(): void { + if (this.closed) return + this.closed = true + this.stream.close() + } +} diff --git a/packages/session-persistence/session-persistence-jsonl/src/zstd-public-decoder.ts b/packages/session-persistence/session-persistence-jsonl/src/zstd-public-decoder.ts new file mode 100644 index 0000000000..bbbb5dd94e --- /dev/null +++ b/packages/session-persistence/session-persistence-jsonl/src/zstd-public-decoder.ts @@ -0,0 +1,38 @@ +/** + * Public-API synchronous Zstandard frame decoder fallback. + * @module dsh-session-persistence-jsonl/zstd-public-decoder + */ + +import { zstdDecompressSync } from 'node:zlib' +import type { ZstdFrameDecoder, ZstdFrameRange } from './zstd.ts' + +/** Multi-frame adapter built exclusively from Node's supported one-shot API. */ +export class PublicZstdFrameDecoder implements ZstdFrameDecoder { + private started = false + private closed = false + + /** @inheritdoc */ + public *decode(source: Buffer, frames: readonly ZstdFrameRange[]): Generator { + if (this.started) throw new Error('Zstandard frame decoder was already started') + if (this.closed) throw new Error('cannot start a closed Zstandard frame decoder') + this.started = true + try { + for (const frame of frames) { + try { + yield zstdDecompressSync(source.subarray(frame.start, frame.end)) + } catch (error) { + throw new Error(`corrupt Zstandard session log: frame at byte ${frame.start} failed validation`, { + cause: error, + }) + } + } + } finally { + this.close() + } + } + + /** @inheritdoc */ + close(): void { + this.closed = true + } +} diff --git a/packages/session-persistence/session-persistence-jsonl/src/zstd.ts b/packages/session-persistence/session-persistence-jsonl/src/zstd.ts index e29747e399..0a4bb9f2ea 100644 --- a/packages/session-persistence/session-persistence-jsonl/src/zstd.ts +++ b/packages/session-persistence/session-persistence-jsonl/src/zstd.ts @@ -5,8 +5,12 @@ * @module dsh-session-persistence-jsonl/zstd */ -import { constants, zstdCompress, zstdDecompress, type ZstdOptions } from 'node:zlib' +import { + constants, zstdCompress, zstdDecompress, zstdDecompressSync, type ZstdOptions, +} from 'node:zlib' import { promisify } from 'node:util' +import { NodePrivateZstdFrameDecoder } from './zstd-private-decoder.ts' +import { PublicZstdFrameDecoder } from './zstd-public-decoder.ts' const ZSTD_MAGIC = 0xFD2FB528 const zstdCompressAsync = promisify(zstdCompress) @@ -117,6 +121,40 @@ export async function decompressZstdFrame(input: Buffer): Promise { return zstdDecompressAsync(input) } +/** + * Synchronously decompress one complete frame and validate its checksum. + * Complete-log readers time-slice repeated calls so the event loop regains + * control without paying one asynchronous native dispatch per frame. + * @param input - one structurally complete Zstandard frame. + * @returns the frame plaintext. + */ +export function decompressZstdFrameSync(input: Buffer): Buffer { + return zstdDecompressSync(input) +} + +/** Common lifecycle for interchangeable synchronous multi-frame decoders. */ +export interface ZstdFrameDecoder { + /** + * Decode and checksum complete frames in source order. Each yielded buffer + * remains valid only until the iterator advances to the next frame. + * @param source - concatenated Zstandard frame bytes. + * @param frames - structurally complete ranges within `source`. + * @returns one plaintext buffer per frame. + */ + decode(source: Buffer, frames: readonly ZstdFrameRange[]): Generator + /** Release decoder-owned resources; repeated calls are harmless. */ + close(): void +} + +/** + * Select the shared private decoder when the running Node 22/24/26 shape is + * compatible, otherwise preserve correctness with the public one-shot API. + * @returns a synchronous decoder with an implementation-independent lifecycle. + */ +export function createZstdFrameDecoder(): ZstdFrameDecoder { + return NodePrivateZstdFrameDecoder.create() ?? new PublicZstdFrameDecoder() +} + /** * Recover available plaintext from a structurally incomplete final frame. * `ZSTD_e_flush` deliberately suppresses final-frame and checksum completion;