File
Blob: src/worker/git/pack/indexer/resolve/reader.ts
| 1 | import { readPackRange } from "@/worker/git/pack/packMeta"; |
| 2 | |
| 3 | import { InflateCursor } from "../inflateCursor"; |
| 4 | import type { PackEntryTable, ResolveOptions } from "../types"; |
| 5 | import { throwIfAborted } from "./errors"; |
| 6 | |
| 7 | export class SequentialReader { |
| 8 | private buf: Uint8Array<ArrayBufferLike> = new Uint8Array(0); |
| 9 | private bufAbsStart = 0; |
| 10 | private env: Env; |
| 11 | private packKey: string; |
| 12 | private packSize: number; |
| 13 | private chunkSize: number; |
| 14 | private limiter: ResolveOptions["limiter"]; |
| 15 | private countSub: ResolveOptions["countSubrequest"]; |
| 16 | private log: ResolveOptions["log"]; |
| 17 | private signal?: AbortSignal; |
| 18 | |
| 19 | constructor( |
| 20 | env: Env, |
| 21 | packKey: string, |
| 22 | packSize: number, |
| 23 | chunkSize: number, |
| 24 | limiter: ResolveOptions["limiter"], |
| 25 | countSub: ResolveOptions["countSubrequest"], |
| 26 | log: ResolveOptions["log"], |
| 27 | signal?: AbortSignal |
| 28 | ) { |
| 29 | this.env = env; |
| 30 | this.packKey = packKey; |
| 31 | this.packSize = packSize; |
| 32 | this.chunkSize = chunkSize; |
| 33 | this.limiter = limiter; |
| 34 | this.countSub = countSub; |
| 35 | this.log = log; |
| 36 | this.signal = signal; |
| 37 | } |
| 38 | |
| 39 | throwIfAborted(stage: string): void { |
| 40 | throwIfAborted(this.signal, this.log, stage); |
| 41 | } |
| 42 | |
| 43 | /** |
| 44 | * Read a byte range from the pack. If the range falls within the current |
| 45 | * buffered chunk, return a subarray. Otherwise preload a new chunk starting |
| 46 | * at the requested offset so nearby follow-on reads stay coalesced. |
| 47 | */ |
| 48 | async readRange(offset: number, length: number): Promise<Uint8Array> { |
| 49 | this.throwIfAborted("reader:read-range"); |
| 50 | const bufEnd = this.bufAbsStart + this.buf.length; |
| 51 | if (offset >= this.bufAbsStart && offset + length <= bufEnd) { |
| 52 | const localStart = offset - this.bufAbsStart; |
| 53 | return this.buf.subarray(localStart, localStart + length); |
| 54 | } |
| 55 | if (length > this.chunkSize) { |
| 56 | const data = await readPackRange(this.env, this.packKey, offset, length, { |
| 57 | limiter: this.limiter, |
| 58 | countSubrequest: this.countSub, |
| 59 | signal: this.signal, |
| 60 | }); |
| 61 | if (!data) { |
| 62 | this.throwIfAborted("reader:read-range"); |
| 63 | throw new Error("resolve: R2 read failure"); |
| 64 | } |
| 65 | return data; |
| 66 | } |
| 67 | |
| 68 | await this.preload(offset); |
| 69 | const preloadEnd = this.bufAbsStart + this.buf.length; |
| 70 | if (offset + length <= preloadEnd) { |
| 71 | const localStart = offset - this.bufAbsStart; |
| 72 | return this.buf.subarray(localStart, localStart + length); |
| 73 | } |
| 74 | |
| 75 | const data = await readPackRange(this.env, this.packKey, offset, length, { |
| 76 | limiter: this.limiter, |
| 77 | countSubrequest: this.countSub, |
| 78 | signal: this.signal, |
| 79 | }); |
| 80 | if (!data) { |
| 81 | this.throwIfAborted("reader:read-range"); |
| 82 | throw new Error("resolve: R2 read failure"); |
| 83 | } |
| 84 | return data; |
| 85 | } |
| 86 | |
| 87 | /** |
| 88 | * Return the largest already-buffered window starting at `offset`, preloading |
| 89 | * a new chunk when needed. Unlike `readRange()`, this intentionally does not |
| 90 | * stitch together the full requested span, so pass-2 inflate can stream large |
| 91 | * entries without double-buffering their compressed bytes. |
| 92 | */ |
| 93 | async readWindow(offset: number, maxLength: number): Promise<Uint8Array> { |
| 94 | this.throwIfAborted("reader:read-window"); |
| 95 | if (maxLength <= 0 || offset >= this.packSize) return new Uint8Array(0); |
| 96 | |
| 97 | const bufEnd = this.bufAbsStart + this.buf.length; |
| 98 | if (!(offset >= this.bufAbsStart && offset < bufEnd)) { |
| 99 | await this.preload(offset); |
| 100 | } |
| 101 | |
| 102 | const windowEnd = this.bufAbsStart + this.buf.length; |
| 103 | if (offset < this.bufAbsStart || offset >= windowEnd) { |
| 104 | const data = await readPackRange(this.env, this.packKey, offset, Math.min(maxLength, 1), { |
| 105 | limiter: this.limiter, |
| 106 | countSubrequest: this.countSub, |
| 107 | signal: this.signal, |
| 108 | }); |
| 109 | if (!data) { |
| 110 | this.throwIfAborted("reader:read-window"); |
| 111 | throw new Error("resolve: R2 read failure"); |
| 112 | } |
| 113 | return data; |
| 114 | } |
| 115 | |
| 116 | const localStart = offset - this.bufAbsStart; |
| 117 | const localLength = Math.min(maxLength, windowEnd - offset); |
| 118 | return this.buf.subarray(localStart, localStart + localLength); |
| 119 | } |
| 120 | |
| 121 | /** Preload a large sequential chunk starting at the given offset. */ |
| 122 | async preload(offset: number): Promise<void> { |
| 123 | this.throwIfAborted("reader:preload"); |
| 124 | const bytesLeft = this.packSize - offset; |
| 125 | if (bytesLeft <= 0) return; |
| 126 | const readLen = Math.min(this.chunkSize, bytesLeft); |
| 127 | const chunk = await readPackRange(this.env, this.packKey, offset, readLen, { |
| 128 | limiter: this.limiter, |
| 129 | countSubrequest: this.countSub, |
| 130 | signal: this.signal, |
| 131 | }); |
| 132 | if (!chunk) { |
| 133 | this.throwIfAborted("reader:preload"); |
| 134 | throw new Error("resolve: R2 preload failure"); |
| 135 | } |
| 136 | this.buf = chunk; |
| 137 | this.bufAbsStart = offset; |
| 138 | } |
| 139 | } |
| 140 | |
| 141 | /** Inflate a pack entry's compressed payload using a buffered pack reader. */ |
| 142 | export async function inflateFromReader( |
| 143 | reader: SequentialReader, |
| 144 | table: PackEntryTable, |
| 145 | index: number |
| 146 | ): Promise<Uint8Array> { |
| 147 | reader.throwIfAborted("reader:inflate-entry"); |
| 148 | const payloadStart = table.offsets[index] + table.headerLens[index]; |
| 149 | const cursor = new InflateCursor(); |
| 150 | let nextOffset = payloadStart; |
| 151 | let firstPush = true; |
| 152 | |
| 153 | while (!cursor.finished) { |
| 154 | reader.throwIfAborted("reader:inflate-entry"); |
| 155 | const bytesLeft = table.spanEnds[index] - nextOffset; |
| 156 | if (bytesLeft <= 0) { |
| 157 | throw new Error(`resolve: incomplete inflate for entry at offset ${table.offsets[index]}`); |
| 158 | } |
| 159 | |
| 160 | const minBytes = firstPush ? 2 : 1; |
| 161 | let window = await reader.readWindow(nextOffset, bytesLeft); |
| 162 | if (window.length < minBytes) { |
| 163 | // A chunk size of 1 is valid in tests and can split the zlib wrapper at |
| 164 | // arbitrary boundaries. Stitch just the minimum prefix needed for the |
| 165 | // inflate cursor to make forward progress. |
| 166 | window = await reader.readRange(nextOffset, Math.min(bytesLeft, minBytes)); |
| 167 | } |
| 168 | if (window.length < minBytes) { |
| 169 | throw new Error( |
| 170 | `resolve: unexpected EOF while inflating entry at offset ${table.offsets[index]}` |
| 171 | ); |
| 172 | } |
| 173 | |
| 174 | cursor.push(window); |
| 175 | firstPush = false; |
| 176 | |
| 177 | const consumed = cursor.consumedInputBytes; |
| 178 | if (consumed <= 0 && !cursor.finished) { |
| 179 | throw new Error(`resolve: inflate stalled at offset ${nextOffset}`); |
| 180 | } |
| 181 | nextOffset += consumed; |
| 182 | } |
| 183 | |
| 184 | if (nextOffset !== table.spanEnds[index]) { |
| 185 | throw new Error( |
| 186 | `resolve: inflate span mismatch at offset ${table.offsets[index]} (expected end ${table.spanEnds[index]}, got ${nextOffset})` |
| 187 | ); |
| 188 | } |
| 189 | return cursor.output; |
| 190 | } |