Skip to content
File

Blob: src/worker/tasks/refBackfill.ts

typescript150 lines
1import type { CacheContext } from "@/worker/cache";
2import type { Logger } from "@/worker/common/logger";
3import type { RepoDurableObject } from "@/worker/do/repo/repoDO";
4 
5import { type PackRefBackfillQueueMessage, type RepoQueueMessageHandle } from "./types";
6 
7import { getRepoStubByDoId } from "@/worker/common";
8import { loadIdxView } from "@/worker/git/object-store";
9import { resolveDeltasAndWriteIdx, scanPack } from "@/worker/git/pack/indexer";
10import { loadPackRefView } from "@/worker/git/pack/refIndex";
11import { createQueueTaskContext, logSoftBudgetExhausted, retryQueueMessage } from "./context";
12 
13const REF_BACKFILL_SUBREQUEST_BUDGET = 7_500;
14const REF_BACKFILL_RETRY_DELAY_SECONDS = 30;
15 
16function countBackfillSubrequest(cacheCtx: CacheContext, log: Logger, op: string, n = 1): void {
17 logSoftBudgetExhausted({
18 cacheCtx,
19 log,
20 flagPrefix: "ref-backfill-soft-budget",
21 op,
22 count: n,
23 });
24}
25 
26function isDeterministicPackFailure(error: unknown): boolean {
27 const message = String(error);
28 return (
29 message.includes("invalid") ||
30 message.includes("mismatch") ||
31 message.includes("unsupported") ||
32 message.includes("truncated") ||
33 message.includes("cannot fit")
34 );
35}
36 
37export async function handlePackRefBackfillMessage(
38 message: Omit<RepoQueueMessageHandle<PackRefBackfillQueueMessage>, "body">,
39 body: PackRefBackfillQueueMessage,
40 env: Env,
41 ctx: ExecutionContext
42): Promise<void> {
43 const repoLabel = body.repoId || `do:${body.doId}`;
44 const task = createQueueTaskContext({
45 env,
46 ctx,
47 repoLabel,
48 operation: "pack-refs",
49 subrequestBudget: REF_BACKFILL_SUBREQUEST_BUDGET,
50 });
51 const log = task.logFor({
52 service: "PackRefBackfillQueue",
53 repoId: repoLabel,
54 doId: body.doId,
55 });
56 const stub = getRepoStubByDoId(env, body.doId) as DurableObjectStub<RepoDurableObject>;
57 const { cacheCtx, limiter } = task;
58 
59 try {
60 log.info("ref-index:backfill-start", { packKey: body.packKey });
61 
62 countBackfillSubrequest(cacheCtx, log, "do:get-active-pack-catalog");
63 const activeCatalog = await limiter.run("do:get-active-pack-catalog", async () => {
64 return await stub.getActivePackCatalog();
65 });
66 cacheCtx.memo = cacheCtx.memo || {};
67 cacheCtx.memo.packCatalog = activeCatalog;
68 
69 const target = activeCatalog.find((row) => row.packKey === body.packKey);
70 if (!target) {
71 log.info("ref-index:backfill-stale-pack", { packKey: body.packKey });
72 message.ack();
73 return;
74 }
75 const externalBaseCatalog = activeCatalog.filter((row) => row.packKey !== target.packKey);
76 log.debug("ref-index:backfill-resolve-catalog", {
77 packKey: target.packKey,
78 activePacks: activeCatalog.length,
79 externalBasePacks: externalBaseCatalog.length,
80 });
81 
82 const idxView = await loadIdxView(env, target.packKey, cacheCtx, target.packBytes);
83 if (!idxView) {
84 log.warn("ref-index:backfill-invalid-pack", {
85 packKey: target.packKey,
86 reason: "missing-or-invalid-idx",
87 });
88 message.ack();
89 return;
90 }
91 
92 const existing = await loadPackRefView(env, target.packKey, idxView, cacheCtx);
93 if (existing.type === "Ready") {
94 log.info("ref-index:backfill-complete", {
95 packKey: target.packKey,
96 result: "already-present",
97 });
98 message.ack();
99 return;
100 }
101 
102 const scanResult = await scanPack({
103 env,
104 packKey: target.packKey,
105 packSize: target.packBytes,
106 limiter,
107 countSubrequest: (n = 1) => countBackfillSubrequest(cacheCtx, log, "r2:scan-pack", n),
108 log,
109 });
110 
111 cacheCtx.memo.packCatalog = externalBaseCatalog;
112 const resolveResult = await resolveDeltasAndWriteIdx({
113 env,
114 packKey: target.packKey,
115 packSize: target.packBytes,
116 limiter,
117 countSubrequest: (n = 1) => countBackfillSubrequest(cacheCtx, log, "r2:resolve-pack", n),
118 log,
119 scanResult,
120 activeCatalog: externalBaseCatalog,
121 cacheCtx,
122 repoId: repoLabel,
123 writeIdx: false,
124 existingIdxView: idxView,
125 });
126 
127 log.info("ref-index:backfill-complete", {
128 packKey: target.packKey,
129 objectCount: resolveResult.objectCount,
130 refIndexBytes: resolveResult.refIndexBytes,
131 });
132 message.ack();
133 } catch (error) {
134 if (isDeterministicPackFailure(error)) {
135 log.warn("ref-index:backfill-invalid-pack", {
136 packKey: body.packKey,
137 error: String(error),
138 });
139 message.ack();
140 return;
141 }
142 
143 log.warn("ref-index:backfill-retry", {
144 packKey: body.packKey,
145 error: String(error),
146 });
147 retryQueueMessage(message, REF_BACKFILL_RETRY_DELAY_SECONDS);
148 }
149}