zstd-private-decoder.ts 6.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178
  1. /**
  2. * Node-private synchronous Zstandard frame decoder optimization.
  3. * @module dsh-session-persistence-jsonl/zstd-private-decoder
  4. */
  5. import { constants as bufferConstants } from 'node:buffer'
  6. import { createZstdDecompress } from 'node:zlib'
  7. import type { ZstdFrameDecoder, ZstdFrameRange } from './zstd.ts'
  8. const DECODE_CHUNK_SIZE = 1024 * 1024
  9. interface NodeZstdPrivateHandle {
  10. writeSync(
  11. flushFlag: number,
  12. input: Buffer,
  13. inputOffset: number,
  14. inputLength: number,
  15. output: Buffer,
  16. outputOffset: number,
  17. outputLength: number,
  18. ): void
  19. }
  20. type NodeZstdPrivateWriteState = Uint32Array & { 0: number; 1: number }
  21. interface NodeZstdPrivateState {
  22. [key: symbol]: unknown
  23. _handle: NodeZstdPrivateHandle | null
  24. _writeState: NodeZstdPrivateWriteState
  25. _defaultFlushFlag: number
  26. }
  27. type NodeZstdPrivateStream = ReturnType<typeof createZstdDecompress> & NodeZstdPrivateState
  28. /** Return the stream with its observed private Node contract, or reject that optimization. */
  29. function privateZstdStream(
  30. stream: ReturnType<typeof createZstdDecompress>,
  31. ): { stream: NodeZstdPrivateStream; errorKey: symbol } | undefined {
  32. const candidate = stream as unknown as Partial<NodeZstdPrivateState>
  33. const handle = candidate._handle
  34. const errorKey = Reflect.ownKeys(stream).find((key): key is symbol => (
  35. typeof key === 'symbol' && key.description === 'kError'
  36. ))
  37. /* v8 ignore next -- one test runtime exposes one Node-private shape; the Node 22/24/26 matrix checks compatibility. */
  38. if (
  39. typeof handle !== 'object' || handle === null
  40. || typeof (handle as { writeSync?: unknown }).writeSync !== 'function'
  41. || !(candidate._writeState instanceof Uint32Array)
  42. || candidate._writeState.length < 2
  43. || typeof candidate._defaultFlushFlag !== 'number'
  44. || errorKey === undefined
  45. || candidate[errorKey] !== null
  46. ) return undefined
  47. return { stream: stream as NodeZstdPrivateStream, errorKey }
  48. }
  49. /**
  50. * Synchronous multi-frame decoder backed by one Node Zstd stream handle. Node
  51. * exposes synchronous decoding only as a one-shot API, so this adapter uses
  52. * the stream's private handle contract to reuse its native context and output
  53. * chunks across frames.
  54. */
  55. export class NodePrivateZstdFrameDecoder implements ZstdFrameDecoder {
  56. private readonly output = Buffer.allocUnsafe(DECODE_CHUNK_SIZE)
  57. private decoderError?: Error
  58. private started = false
  59. private closed = false
  60. private constructor(
  61. private readonly stream: NodeZstdPrivateStream,
  62. private readonly errorKey: symbol,
  63. ) {
  64. this.stream.on('error', (error: Error) => {
  65. this.decoderError ??= error
  66. })
  67. }
  68. /**
  69. * Create the optimized decoder when this Node release exposes the expected
  70. * private stream shape.
  71. * @returns a shared decoder, or `undefined` when callers must use the public fallback.
  72. */
  73. static create(): NodePrivateZstdFrameDecoder | undefined {
  74. const stream = createZstdDecompress({ chunkSize: DECODE_CHUNK_SIZE })
  75. const privateAccess = privateZstdStream(stream)
  76. /* v8 ignore next -- reached only when a supported Node release changes its private stream shape. */
  77. if (privateAccess !== undefined) {
  78. return new NodePrivateZstdFrameDecoder(privateAccess.stream, privateAccess.errorKey)
  79. }
  80. /* v8 ignore next -- the active Node runtime passed the private-shape probe above. */
  81. stream.close()
  82. /* v8 ignore next -- the active Node runtime passed the private-shape probe above. */
  83. return undefined
  84. }
  85. /** @inheritdoc */
  86. public *decode(source: Buffer, frames: readonly ZstdFrameRange[]): Generator<Buffer, void, void> {
  87. if (this.started) throw new Error('Zstandard frame decoder was already started')
  88. if (this.closed) throw new Error('cannot start a closed Zstandard frame decoder')
  89. this.started = true
  90. try {
  91. for (const frame of frames) {
  92. try {
  93. yield this.decodeFrame(source.subarray(frame.start, frame.end))
  94. } catch (error) {
  95. throw new Error(`corrupt Zstandard session log: frame at byte ${frame.start} failed validation`, {
  96. cause: error,
  97. })
  98. }
  99. }
  100. } finally {
  101. this.close()
  102. }
  103. }
  104. /** Decode one frame; its returned scratch view remains valid until the next call. */
  105. private decodeFrame(input: Buffer): Buffer {
  106. const handle = this.stream._handle
  107. /* v8 ignore next -- decode() rejects closed instances before entering this private frame operation. */
  108. if (this.closed || handle === null) throw new Error('cannot decode with a closed Zstandard frame decoder')
  109. let inputOffset = 0
  110. let inputRemaining = input.length
  111. let outputBytes = 0
  112. const fullChunks: Buffer[] = []
  113. for (;;) {
  114. handle.writeSync(
  115. this.stream._defaultFlushFlag,
  116. input,
  117. inputOffset,
  118. inputRemaining,
  119. this.output,
  120. 0,
  121. this.output.length,
  122. )
  123. if (this.decoderError !== undefined) throw this.decoderError
  124. const internalError = this.stream[this.errorKey]
  125. if (internalError !== null) {
  126. if (internalError instanceof Error) throw internalError
  127. throw new Error('Zstandard decoder exposed a non-Error internal failure')
  128. }
  129. const outputAfter = this.stream._writeState[0]
  130. const inputAfter = this.stream._writeState[1]
  131. const consumed = inputRemaining - inputAfter
  132. const produced = this.output.length - outputAfter
  133. if (produced > 0) {
  134. outputBytes += produced
  135. /* v8 ignore next -- Buffer cannot materialize a frame beyond its own process-wide maximum length. */
  136. if (outputBytes > bufferConstants.MAX_LENGTH) {
  137. throw new Error(`Zstandard frame output exceeds ${bufferConstants.MAX_LENGTH} bytes`)
  138. }
  139. }
  140. if (outputAfter !== 0) {
  141. /* v8 ignore next -- structurally scanned ranges contain exactly one complete frame and no trailing bytes. */
  142. if (inputAfter !== 0) throw new Error('Zstandard frame decoder left trailing input')
  143. const finalChunk = this.output.subarray(0, produced)
  144. if (fullChunks.length === 0) return finalChunk
  145. if (produced > 0) fullChunks.push(Buffer.from(finalChunk))
  146. const onlyChunk = fullChunks[0] as Buffer
  147. return fullChunks.length === 1
  148. ? onlyChunk
  149. : Buffer.concat(fullChunks, outputBytes)
  150. }
  151. fullChunks.push(Buffer.from(this.output))
  152. inputOffset += consumed
  153. inputRemaining = inputAfter
  154. }
  155. }
  156. /** @inheritdoc */
  157. close(): void {
  158. if (this.closed) return
  159. this.closed = true
  160. this.stream.close()
  161. }
  162. }