File
Blob: src/worker/git/operations/fetch/neededFast.ts
| 1 | import type { CacheContext } from "@/worker/cache"; |
| 2 | |
| 3 | import { createLogger } from "@/worker/common"; |
| 4 | import { parseCommitRefs } from "@/worker/git/core"; |
| 5 | import { readObject, readObjectRefsBatch } from "@/worker/git/object-store"; |
| 6 | import { findCommonHaves } from "../closure"; |
| 7 | |
| 8 | export async function computeNeededFast( |
| 9 | env: Env, |
| 10 | repoId: string, |
| 11 | wants: string[], |
| 12 | haves: string[], |
| 13 | cacheCtx?: CacheContext, |
| 14 | onProgress?: (message: string) => void |
| 15 | ): Promise<string[]> { |
| 16 | const log = createLogger(env.LOG_LEVEL, { service: "NeededFast", repoId }); |
| 17 | const startTime = Date.now(); |
| 18 | const timeoutMs = 49_000; |
| 19 | |
| 20 | log.debug("fast:building-stop-set", { haves: haves.length }); |
| 21 | const stopSet = new Set<string>(); |
| 22 | |
| 23 | let ackOids: string[] = []; |
| 24 | if (haves.length > 0) { |
| 25 | onProgress?.("Finding common commits...\n"); |
| 26 | ackOids = await findCommonHaves(env, repoId, haves, cacheCtx); |
| 27 | for (const oid of ackOids) { |
| 28 | stopSet.add(oid.toLowerCase()); |
| 29 | } |
| 30 | |
| 31 | if (ackOids.length === 0) { |
| 32 | log.debug("fast:no-common-base", { haves: haves.length }); |
| 33 | } |
| 34 | } |
| 35 | |
| 36 | onProgress?.("Selecting objects to send...\n"); |
| 37 | |
| 38 | if (ackOids.length > 0 && ackOids.length < 10) { |
| 39 | const mainlineBudget = 20; |
| 40 | const mainlineQueue = [...ackOids]; |
| 41 | let walked = 0; |
| 42 | |
| 43 | while (mainlineQueue.length > 0 && walked < mainlineBudget) { |
| 44 | if (Date.now() - startTime > 2_000) break; |
| 45 | |
| 46 | const oid = mainlineQueue.shift()!; |
| 47 | const object = await readObject(env, repoId, oid, cacheCtx); |
| 48 | if (object?.type !== "commit") continue; |
| 49 | |
| 50 | const refs = parseCommitRefs(object.payload); |
| 51 | const parent = refs.parents[0]; |
| 52 | if (!parent || stopSet.has(parent)) continue; |
| 53 | |
| 54 | stopSet.add(parent); |
| 55 | mainlineQueue.push(parent); |
| 56 | walked++; |
| 57 | } |
| 58 | |
| 59 | log.debug("fast:mainline-enriched", { stopSize: stopSet.size, walked }); |
| 60 | } |
| 61 | |
| 62 | const seen = new Set<string>(); |
| 63 | const needed = new Set<string>(); |
| 64 | const queue = [...wants]; |
| 65 | |
| 66 | if (cacheCtx) { |
| 67 | cacheCtx.memo = cacheCtx.memo || {}; |
| 68 | cacheCtx.memo.refs = cacheCtx.memo.refs || new Map<string, string[]>(); |
| 69 | cacheCtx.memo.flags = cacheCtx.memo.flags || new Set<string>(); |
| 70 | } |
| 71 | |
| 72 | let refsBatchCalls = 0; |
| 73 | let memoRefsHits = 0; |
| 74 | let missingRefs = 0; |
| 75 | |
| 76 | log.info("fast:starting-closure", { wants: wants.length, stopSet: stopSet.size }); |
| 77 | |
| 78 | while (queue.length > 0) { |
| 79 | if (Date.now() - startTime > timeoutMs) { |
| 80 | log.warn("fast:timeout", { seen: seen.size, needed: needed.size }); |
| 81 | if (cacheCtx) { |
| 82 | cacheCtx.memo = cacheCtx.memo || {}; |
| 83 | cacheCtx.memo.flags = cacheCtx.memo.flags || new Set<string>(); |
| 84 | cacheCtx.memo.flags.add("closure-timeout"); |
| 85 | } |
| 86 | break; |
| 87 | } |
| 88 | |
| 89 | const batch = queue.splice(0, Math.min(128, queue.length)); |
| 90 | const unseenBatch = batch.filter((oid) => !seen.has(oid)); |
| 91 | if (unseenBatch.length === 0) continue; |
| 92 | |
| 93 | const toProcess: string[] = []; |
| 94 | for (const oid of unseenBatch) { |
| 95 | seen.add(oid); |
| 96 | const oidLc = oid.toLowerCase(); |
| 97 | if (stopSet.has(oidLc)) { |
| 98 | log.debug("fast:hit-stop", { oid }); |
| 99 | continue; |
| 100 | } |
| 101 | |
| 102 | needed.add(oid); |
| 103 | toProcess.push(oid); |
| 104 | } |
| 105 | |
| 106 | if (toProcess.length === 0) continue; |
| 107 | |
| 108 | const refsMap = new Map<string, string[]>(); |
| 109 | if (cacheCtx?.memo?.refs) { |
| 110 | for (const oid of toProcess) { |
| 111 | const refs = cacheCtx.memo.refs.get(oid.toLowerCase()); |
| 112 | if (refs === undefined) continue; |
| 113 | refsMap.set(oid, refs); |
| 114 | memoRefsHits++; |
| 115 | } |
| 116 | } |
| 117 | |
| 118 | const batchOids = toProcess.filter((oid) => !refsMap.has(oid)); |
| 119 | if (batchOids.length > 0) { |
| 120 | try { |
| 121 | const batchMap = await readObjectRefsBatch(env, repoId, batchOids, cacheCtx); |
| 122 | refsBatchCalls++; |
| 123 | |
| 124 | for (const oid of batchOids) { |
| 125 | const refs = batchMap.get(oid); |
| 126 | if (refs === undefined) continue; |
| 127 | |
| 128 | refsMap.set(oid, refs); |
| 129 | if (cacheCtx?.memo) { |
| 130 | cacheCtx.memo.refs = cacheCtx.memo.refs || new Map<string, string[]>(); |
| 131 | cacheCtx.memo.refs.set(oid.toLowerCase(), refs); |
| 132 | } |
| 133 | } |
| 134 | } catch (error) { |
| 135 | log.debug("fast:batch-error", { error: String(error) }); |
| 136 | } |
| 137 | } |
| 138 | |
| 139 | missingRefs += toProcess.length - refsMap.size; |
| 140 | for (const refs of refsMap.values()) { |
| 141 | for (const ref of refs) { |
| 142 | if (!seen.has(ref)) { |
| 143 | queue.push(ref); |
| 144 | } |
| 145 | } |
| 146 | } |
| 147 | } |
| 148 | |
| 149 | log.info("fast:completed", { |
| 150 | needed: needed.size, |
| 151 | seen: seen.size, |
| 152 | stopSet: stopSet.size, |
| 153 | memoHits: memoRefsHits, |
| 154 | refsBatches: refsBatchCalls, |
| 155 | missingRefs, |
| 156 | timeMs: Date.now() - startTime, |
| 157 | }); |
| 158 | |
| 159 | return Array.from(needed); |
| 160 | } |