Skip to content
File

Blob: src/worker/git/pack/indexer/resolve/reader.ts

typescript191 lines
1import { readPackRange } from "@/worker/git/pack/packMeta";
2 
3import { InflateCursor } from "../inflateCursor";
4import type { PackEntryTable, ResolveOptions } from "../types";
5import { throwIfAborted } from "./errors";
6 
7export 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. */
142export 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}