Skip to content
File

Blob: src/worker/tasks/compaction.ts

typescript419 lines
1import type { CacheContext } from "@/worker/cache";
2import type { Logger } from "@/worker/common/logger";
3import type { OrderedPackSnapshot } from "@/worker/git/operations/fetch/types";
4import type { RepoDurableObject } from "@/worker/do/repo/repoDO";
5 
6import {
7 type CompactionDeleteQueueMessage,
8 type CompactionQueueMessage,
9 type RepoQueueMessageHandle,
10} from "./types";
11 
12import { getRepoStubByDoId } from "@/worker/common";
13import { buildCompactionNeededOids } from "@/worker/git/compaction/plan";
14import { type SubrequestLimiter } from "@/worker/git/operations/limits";
15import { scanPack, resolveDeltasAndWriteIdx } from "@/worker/git/pack/indexer";
16import { rewritePackResult } from "@/worker/git/pack/rewrite";
17import { loadOrderedPackSnapshot } from "@/worker/git/pack/snapshot";
18import {
19 deleteStagedPack,
20 stagePackToR2,
21 type StagedPackUpload,
22} from "@/worker/git/receive/r2Upload";
23import { doPrefix, packIndexKey, packRefsKey, r2PackKey } from "@/worker/keys";
24import { createQueueTaskContext, logSoftBudgetExhausted, retryQueueMessage } from "./context";
25 
26const COMPACTION_SUBREQUEST_BUDGET = 7_500;
27const COMPACTION_RETRY_DELAY_SECONDS = 30;
28const COMPACTION_CONFLICT_RETRY_DELAY_SECONDS = 10;
29const COMPACTION_DELETE_DELAY_SECONDS = 60;
30 
31function countCompactionSubrequest(cacheCtx: CacheContext, log: Logger, op: string, n = 1): void {
32 logSoftBudgetExhausted({
33 cacheCtx,
34 log,
35 flagPrefix: "compaction-soft-budget",
36 op,
37 count: n,
38 });
39}
40 
41async function cleanupStagedCompaction(args: {
42 stagedUpload: StagedPackUpload | undefined;
43 log: Logger;
44 reason: string;
45}) {
46 if (!args.stagedUpload) return;
47 try {
48 await deleteStagedPack(args.stagedUpload);
49 } catch (error) {
50 args.log.warn("compaction:cleanup-failed", {
51 reason: args.reason,
52 packKey: args.stagedUpload.packKey,
53 error: String(error),
54 });
55 }
56}
57 
58async function abortCompactionLease(args: {
59 stub: DurableObjectStub<RepoDurableObject>;
60 leaseToken: string | undefined;
61 limiter: SubrequestLimiter;
62 cacheCtx: CacheContext;
63 log: Logger;
64 reason: string;
65}) {
66 const leaseToken = args.leaseToken;
67 if (!leaseToken) return;
68 try {
69 countCompactionSubrequest(args.cacheCtx, args.log, "do:abort-compaction");
70 const cleared = await args.limiter.run("do:abort-compaction", async () => {
71 return await args.stub.abortCompaction(leaseToken);
72 });
73 if (!cleared) {
74 args.log.warn("compaction:abort-missed", {
75 reason: args.reason,
76 leaseToken: args.leaseToken,
77 });
78 return;
79 }
80 args.log.info("compaction:abort-complete", {
81 reason: args.reason,
82 leaseToken,
83 });
84 } catch (error) {
85 args.log.warn("compaction:abort-failed", {
86 reason: args.reason,
87 leaseToken: args.leaseToken,
88 error: String(error),
89 });
90 }
91}
92 
93async function clearCompactionRequestAfterBlocked(args: {
94 stub: DurableObjectStub<RepoDurableObject>;
95 limiter: SubrequestLimiter;
96 cacheCtx: CacheContext;
97 log: Logger;
98 reason: string;
99}): Promise<void> {
100 try {
101 countCompactionSubrequest(args.cacheCtx, args.log, "do:clear-compaction-request");
102 await args.limiter.run("do:clear-compaction-request", async () => {
103 await args.stub.clearCompactionRequest();
104 });
105 args.log.warn("compaction:blocked-cleared", { reason: args.reason });
106 } catch (error) {
107 args.log.warn("compaction:blocked-clear-failed", {
108 reason: args.reason,
109 error: String(error),
110 });
111 }
112}
113 
114export async function handleCompactionMessage(
115 message: Omit<RepoQueueMessageHandle<CompactionQueueMessage>, "body">,
116 body: CompactionQueueMessage,
117 env: Env,
118 ctx: ExecutionContext
119): Promise<void> {
120 const repoLabel = body.repoId || `do:${body.doId}`;
121 const task = createQueueTaskContext({
122 env,
123 ctx,
124 repoLabel,
125 operation: "compaction",
126 subrequestBudget: COMPACTION_SUBREQUEST_BUDGET,
127 });
128 const log = task.logFor({
129 service: "CompactionQueue",
130 repoId: repoLabel,
131 doId: body.doId,
132 });
133 const stub = getRepoStubByDoId(env, body.doId) as DurableObjectStub<RepoDurableObject>;
134 const { cacheCtx, limiter } = task;
135 
136 let stagedUpload: StagedPackUpload | undefined;
137 let leaseToken: string | undefined;
138 
139 try {
140 countCompactionSubrequest(cacheCtx, log, "do:begin-compaction");
141 const begin = await limiter.run("do:begin-compaction", async () => {
142 return await stub.beginCompaction();
143 });
144 if (!begin.ok) {
145 if (begin.status === "busy" && begin.reason === "receive-active") {
146 log.info("compaction:busy-retry", { reason: begin.reason });
147 retryQueueMessage(message, COMPACTION_CONFLICT_RETRY_DELAY_SECONDS);
148 return;
149 }
150 
151 log.info("compaction:skip", {
152 status: begin.status,
153 reason: "reason" in begin ? begin.reason : undefined,
154 });
155 message.ack();
156 return;
157 }
158 
159 leaseToken = begin.lease.token;
160 cacheCtx.memo = cacheCtx.memo || {};
161 cacheCtx.memo.packCatalog = begin.activeCatalog;
162 
163 const snapshotLoad = await loadOrderedPackSnapshot(env, repoLabel, cacheCtx, log);
164 if (snapshotLoad.type !== "Ready") {
165 log.warn("compaction:snapshot-unavailable", { reason: snapshotLoad.reason });
166 await abortCompactionLease({
167 stub,
168 leaseToken,
169 limiter,
170 cacheCtx,
171 log,
172 reason: snapshotLoad.reason,
173 });
174 retryQueueMessage(message, COMPACTION_RETRY_DELAY_SECONDS);
175 return;
176 }
177 
178 const snapshot = snapshotLoad.snapshot;
179 const sourcePackMap = new Map(snapshot.packs.map((pack) => [pack.packKey, pack]));
180 const sourcePacks = begin.sourcePacks
181 .map((row) => sourcePackMap.get(row.packKey))
182 .filter((pack): pack is (typeof snapshot.packs)[number] => pack !== undefined);
183 if (sourcePacks.length !== begin.sourcePacks.length) {
184 log.warn("compaction:source-pack-missing", {
185 expected: begin.sourcePacks.length,
186 actual: sourcePacks.length,
187 });
188 await abortCompactionLease({
189 stub,
190 leaseToken,
191 limiter,
192 cacheCtx,
193 log,
194 reason: "source-pack-missing",
195 });
196 retryQueueMessage(message, COMPACTION_CONFLICT_RETRY_DELAY_SECONDS);
197 return;
198 }
199 
200 const neededOids = buildCompactionNeededOids(sourcePacks);
201 log.info("compaction:rewrite-start", {
202 sourceTier: begin.targetTier - 1,
203 targetTier: begin.targetTier,
204 sourceCount: begin.sourcePacks.length,
205 neededCount: neededOids.length,
206 });
207 
208 // Build a compaction-specific snapshot: source packs first so
209 // resolveOrderedEntryByOid picks authoritative source entries for needed
210 // OIDs, then remaining active packs in their normal newest-first order
211 // for delta base closure. Without this reorder, a duplicate identity
212 // REF_DELTA in a newer non-source pack can shadow the source entry and
213 // create a self-referential delta cycle in the topology sort.
214 const sourceKeySet = new Set(begin.sourcePacks.map((row) => row.packKey));
215 const fallbackPacks = snapshot.packs.filter((pack) => !sourceKeySet.has(pack.packKey));
216 const compactionSnapshot: OrderedPackSnapshot = {
217 packs: [...sourcePacks, ...fallbackPacks],
218 };
219 
220 const rewriteResult = await rewritePackResult(env, compactionSnapshot, neededOids, {
221 limiter,
222 countSubrequest: (n) => countCompactionSubrequest(cacheCtx, log, "r2:rewrite-pack", n),
223 });
224 if (rewriteResult.status !== "ok") {
225 log.warn("compaction:rewrite-unavailable", {
226 reason: rewriteResult.failure.reason,
227 retryable: rewriteResult.failure.retryable,
228 details: rewriteResult.failure.details,
229 });
230 await abortCompactionLease({
231 stub,
232 leaseToken,
233 limiter,
234 cacheCtx,
235 log,
236 reason: rewriteResult.failure.reason,
237 });
238 
239 if (rewriteResult.failure.retryable) {
240 retryQueueMessage(message, COMPACTION_RETRY_DELAY_SECONDS);
241 return;
242 }
243 
244 leaseToken = undefined;
245 await clearCompactionRequestAfterBlocked({
246 stub,
247 limiter,
248 cacheCtx,
249 log,
250 reason: rewriteResult.failure.reason,
251 });
252 log.error("compaction:blocked", {
253 reason: rewriteResult.failure.reason,
254 sourceCount: begin.sourcePacks.length,
255 sourceSeqLo: begin.sourcePacks[0]?.seqLo,
256 sourceSeqHi: begin.sourcePacks[begin.sourcePacks.length - 1]?.seqHi,
257 });
258 message.ack();
259 return;
260 }
261 
262 const packKey = r2PackKey(doPrefix(body.doId), `pack-cmp-${begin.lease.token}.pack`);
263 stagedUpload = await stagePackToR2({
264 env,
265 request: new Request(`https://queue.internal/${encodeURIComponent(repoLabel)}/compact-pack`),
266 packStream: rewriteResult.stream,
267 packKey,
268 bytesConsumed: 0,
269 limiter,
270 countSubrequest: (op, n = 1) => countCompactionSubrequest(cacheCtx, log, op, n),
271 });
272 
273 const scanResult = await scanPack({
274 env,
275 packKey: stagedUpload.packKey,
276 packSize: stagedUpload.packBytes,
277 limiter,
278 countSubrequest: (n = 1) => countCompactionSubrequest(cacheCtx, log, "r2:scan-pack", n),
279 log,
280 });
281 const resolveResult = await resolveDeltasAndWriteIdx({
282 env,
283 packKey: stagedUpload.packKey,
284 packSize: stagedUpload.packBytes,
285 limiter,
286 countSubrequest: (n = 1) => countCompactionSubrequest(cacheCtx, log, "r2:resolve-pack", n),
287 log,
288 scanResult,
289 activeCatalog: begin.activeCatalog,
290 cacheCtx,
291 repoId: repoLabel,
292 });
293 
294 countCompactionSubrequest(cacheCtx, log, "do:commit-compaction");
295 const committedUpload = stagedUpload;
296 const commit = await limiter.run("do:commit-compaction", async () => {
297 return await stub.commitCompaction({
298 token: begin.lease.token,
299 sourcePacks: begin.sourcePacks,
300 targetTier: begin.targetTier,
301 packsetVersion: begin.packsetVersion,
302 stagedPack: {
303 packKey: committedUpload.packKey,
304 packBytes: committedUpload.packBytes,
305 idxBytes: resolveResult.idxBytes,
306 objectCount: resolveResult.objectCount,
307 },
308 });
309 });
310 
311 if (commit.status === "retry") {
312 await cleanupStagedCompaction({
313 stagedUpload,
314 log,
315 reason: commit.reason,
316 });
317 leaseToken = undefined;
318 log.info("compaction:retry", { reason: commit.reason });
319 retryQueueMessage(message, COMPACTION_CONFLICT_RETRY_DELAY_SECONDS);
320 return;
321 }
322 
323 leaseToken = undefined;
324 if (commit.shouldRequeue) {
325 ctx.waitUntil(
326 env.REPO_TASKS_QUEUE.send({
327 kind: "compaction",
328 doId: body.doId,
329 repoId: body.repoId,
330 }).catch((error) => {
331 log.warn("compaction:follow-up-enqueue-failed", { error: String(error) });
332 })
333 );
334 }
335 
336 if (commit.supersededPackKeys.length > 0) {
337 ctx.waitUntil(
338 env.REPO_TASKS_QUEUE.send(
339 {
340 kind: "compaction-delete",
341 doId: body.doId,
342 repoId: body.repoId,
343 packKeys: commit.supersededPackKeys,
344 },
345 { delaySeconds: COMPACTION_DELETE_DELAY_SECONDS }
346 ).catch((error) => {
347 log.warn("compaction:delete-enqueue-failed", { error: String(error) });
348 })
349 );
350 }
351 
352 log.info("compaction:done", {
353 targetPackKey: commit.targetPackKey,
354 supersededCount: commit.supersededPackKeys.length,
355 shouldRequeue: commit.shouldRequeue,
356 });
357 message.ack();
358 } catch (error) {
359 log.error("compaction:error", { error: String(error) });
360 await cleanupStagedCompaction({
361 stagedUpload,
362 log,
363 reason: "error",
364 });
365 await abortCompactionLease({
366 stub,
367 leaseToken,
368 limiter,
369 cacheCtx,
370 log,
371 reason: "error",
372 });
373 retryQueueMessage(message, COMPACTION_RETRY_DELAY_SECONDS);
374 }
375}
376 
377export async function handleCompactionDeleteMessage(
378 message: Omit<RepoQueueMessageHandle<CompactionDeleteQueueMessage>, "body">,
379 body: CompactionDeleteQueueMessage,
380 env: Env,
381 ctx: ExecutionContext
382): Promise<void> {
383 const repoLabel = body.repoId || `do:${body.doId}`;
384 const task = createQueueTaskContext({
385 env,
386 ctx,
387 repoLabel,
388 operation: "compaction-delete",
389 subrequestBudget: 25,
390 });
391 const log = task.logFor({
392 service: "CompactionDeleteQueue",
393 repoId: repoLabel,
394 doId: body.doId,
395 });
396 const { limiter } = task;
397 
398 try {
399 const keysToDelete: string[] = [];
400 for (const packKey of body.packKeys) {
401 // Each superseded pack has three derived immutable artifacts in R2:
402 // the pack bytes, the idx, and the logical-reference sidecar.
403 keysToDelete.push(packKey, packIndexKey(packKey), packRefsKey(packKey));
404 }
405 
406 await limiter.run("r2:delete-superseded-packs", async () => {
407 await env.REPO_BUCKET.delete(keysToDelete);
408 });
409 log.info("compaction:delete-complete", {
410 packCount: body.packKeys.length,
411 artifactCount: keysToDelete.length,
412 });
413 message.ack();
414 } catch (error) {
415 log.warn("compaction:delete-failed", { error: String(error) });
416 retryQueueMessage(message, COMPACTION_RETRY_DELAY_SECONDS);
417 }
418}