Skip to content
File

Blob: src/worker/git/object-store/store.ts

typescript197 lines
1import type { CacheContext } from "@/worker/cache";
2import type { Logger } from "@/worker/common/logger";
3import type { PackedObjectResult } from "./types";
4 
5import { createBlobFromBytes } from "@/worker/common";
6import { parseCommitRefs, parseTagTarget, parseTreeChildOids } from "@/worker/git/core";
7import {
8 MAX_SIMULTANEOUS_CONNECTIONS,
9 countSubrequest,
10 getLimiter,
11} from "@/worker/git/operations/limits";
12import { findObject } from "./lookup";
13import { materializePackedObjectCandidate } from "./materialize";
14import { ensureMemo, getPackedObjectStoreLogger, logOnce, type ResolvedLocation } from "./support";
15 
16function countPackedSubrequest(
17 cacheCtx: CacheContext | undefined,
18 log: Logger,
19 details: { op: string; oid?: string; packKey?: string },
20 flag: string,
21 n?: number
22) {
23 if (countSubrequest(cacheCtx, n)) return;
24 logOnce(cacheCtx, flag, () => {
25 log.warn("soft-budget-exhausted", details);
26 });
27}
28 
29async function readObjectFromLocation(
30 env: Env,
31 repoId: string,
32 location: ResolvedLocation,
33 cacheCtx: CacheContext | undefined,
34 visited: Set<string>
35): Promise<PackedObjectResult | undefined> {
36 const limiter = getLimiter(cacheCtx);
37 const log = getPackedObjectStoreLogger(env, repoId);
38 
39 // Object-store reads intentionally keep first-hit REF_DELTA semantics. The
40 // indexer backfill path is the only caller that tries alternate duplicates.
41 return await materializePackedObjectCandidate({
42 env,
43 candidate: location,
44 limiter,
45 countSubrequest: (n?: number) => {
46 countPackedSubrequest(
47 cacheCtx,
48 log,
49 {
50 op: "r2:get-pack-entry",
51 oid: location.oid,
52 packKey: location.source.packKey,
53 },
54 "packed-read-entry-soft-budget-warned",
55 n
56 );
57 },
58 log,
59 cyclePolicy: "throw",
60 resolveRefBase: async (baseOid, nextVisited) => {
61 return await readObject(env, repoId, baseOid, cacheCtx, nextVisited);
62 },
63 visited,
64 });
65}
66 
67export async function readObject(
68 env: Env,
69 repoId: string,
70 oid: string,
71 cacheCtx?: CacheContext,
72 visited?: Set<string>
73): Promise<PackedObjectResult | undefined> {
74 const oidLc = oid.toLowerCase();
75 ensureMemo(cacheCtx, repoId);
76 const log = getPackedObjectStoreLogger(env, repoId);
77 
78 const cached = cacheCtx?.memo?.packedObjects?.get(oidLc);
79 if (cached !== undefined) return cached || undefined;
80 
81 const inflight = cacheCtx?.memo?.packedObjectPromises?.get(oidLc);
82 if (inflight) return await inflight;
83 
84 const promise = (async () => {
85 const location = await findObject(env, repoId, oidLc, cacheCtx);
86 if (!location) return undefined;
87 return await readObjectFromLocation(env, repoId, location, cacheCtx, visited || new Set());
88 })();
89 
90 if (cacheCtx?.memo) {
91 cacheCtx.memo.packedObjectPromises = cacheCtx.memo.packedObjectPromises || new Map();
92 cacheCtx.memo.packedObjectPromises.set(oidLc, promise);
93 }
94 
95 try {
96 const result = await promise;
97 if (cacheCtx?.memo) {
98 cacheCtx.memo.packedObjects = cacheCtx.memo.packedObjects || new Map();
99 cacheCtx.memo.packedObjects.set(oidLc, result || null);
100 }
101 if (result) {
102 logOnce(cacheCtx, "packed-object-read-logged", () => {
103 log.debug("object-read", {
104 source: "pack-catalog",
105 packKey: result.packKey,
106 type: result.type,
107 });
108 });
109 }
110 return result;
111 } finally {
112 cacheCtx?.memo?.packedObjectPromises?.delete(oidLc);
113 }
114}
115 
116export async function hasObjectsBatch(
117 env: Env,
118 repoId: string,
119 oids: string[],
120 cacheCtx?: CacheContext
121): Promise<boolean[]> {
122 const results: boolean[] = [];
123 for (let i = 0; i < oids.length; i += MAX_SIMULTANEOUS_CONNECTIONS) {
124 const batch = oids.slice(i, i + MAX_SIMULTANEOUS_CONNECTIONS);
125 const batchResults = await Promise.all(
126 batch.map(async (oid) => {
127 const found = await findObject(env, repoId, oid, cacheCtx);
128 return !!found;
129 })
130 );
131 results.push(...batchResults);
132 }
133 return results;
134}
135 
136export async function readObjectRefsBatch(
137 env: Env,
138 repoId: string,
139 oids: string[],
140 cacheCtx?: CacheContext
141): Promise<Map<string, string[]>> {
142 const out = new Map<string, string[]>();
143 for (let index = 0; index < oids.length; index += MAX_SIMULTANEOUS_CONNECTIONS) {
144 const batch = oids.slice(index, index + MAX_SIMULTANEOUS_CONNECTIONS);
145 const objects = await Promise.all(batch.map((oid) => readObject(env, repoId, oid, cacheCtx)));
146 
147 for (let batchIndex = 0; batchIndex < batch.length; batchIndex++) {
148 const oid = batch[batchIndex];
149 const obj = objects[batchIndex];
150 if (!obj) {
151 // Omit missing objects so fetch closure can return the partial pack-first
152 // result it actually discovered instead of inventing compatibility reads.
153 continue;
154 }
155 if (obj.type === "commit") {
156 const refs = parseCommitRefs(obj.payload);
157 out.set(
158 oid,
159 [refs.tree, ...refs.parents].filter((value): value is string => !!value)
160 );
161 continue;
162 }
163 if (obj.type === "tree") {
164 out.set(oid, parseTreeChildOids(obj.payload));
165 continue;
166 }
167 if (obj.type === "tag") {
168 const tag = parseTagTarget(obj.payload);
169 out.set(oid, tag?.targetOid ? [tag.targetOid] : []);
170 continue;
171 }
172 out.set(oid, []);
173 }
174 }
175 return out;
176}
177 
178export async function readBlobStream(
179 env: Env,
180 repoId: string,
181 oid: string,
182 cacheCtx?: CacheContext
183): Promise<Response | null> {
184 const obj = await readObject(env, repoId, oid, cacheCtx);
185 if (!obj || obj.type !== "blob") return null;
186 return new Response(createBlobFromBytes(obj.payload).stream(), {
187 headers: {
188 "Content-Type": "application/octet-stream",
189 "Cache-Control": "public, max-age=31536000, immutable",
190 ETag: `"${obj.oid}"`,
191 },
192 });
193}
194 
195export { findObject } from "./lookup";
196export { logPackedObjectMismatch } from "./support";