File
Blob: src/worker/git/object-store/store.ts
| 1 | import type { CacheContext } from "@/worker/cache"; |
| 2 | import type { Logger } from "@/worker/common/logger"; |
| 3 | import type { PackedObjectResult } from "./types"; |
| 4 | |
| 5 | import { createBlobFromBytes } from "@/worker/common"; |
| 6 | import { parseCommitRefs, parseTagTarget, parseTreeChildOids } from "@/worker/git/core"; |
| 7 | import { |
| 8 | MAX_SIMULTANEOUS_CONNECTIONS, |
| 9 | countSubrequest, |
| 10 | getLimiter, |
| 11 | } from "@/worker/git/operations/limits"; |
| 12 | import { findObject } from "./lookup"; |
| 13 | import { materializePackedObjectCandidate } from "./materialize"; |
| 14 | import { ensureMemo, getPackedObjectStoreLogger, logOnce, type ResolvedLocation } from "./support"; |
| 15 | |
| 16 | function 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 | |
| 29 | async 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 | |
| 67 | export 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 | |
| 116 | export 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 | |
| 136 | export 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 | |
| 178 | export 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 | |
| 195 | export { findObject } from "./lookup"; |
| 196 | export { logPackedObjectMismatch } from "./support"; |