Skip to content
File

Blob: src/worker/git/operations/fetch/neededFast.ts

typescript161 lines
1import type { CacheContext } from "@/worker/cache";
2 
3import { createLogger } from "@/worker/common";
4import { parseCommitRefs } from "@/worker/git/core";
5import { readObject, readObjectRefsBatch } from "@/worker/git/object-store";
6import { findCommonHaves } from "../closure";
7 
8export 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}