File
Blob: src/worker/tasks/refBackfill.ts
| 1 | import type { CacheContext } from "@/worker/cache"; |
| 2 | import type { Logger } from "@/worker/common/logger"; |
| 3 | import type { RepoDurableObject } from "@/worker/do/repo/repoDO"; |
| 4 | |
| 5 | import { type PackRefBackfillQueueMessage, type RepoQueueMessageHandle } from "./types"; |
| 6 | |
| 7 | import { getRepoStubByDoId } from "@/worker/common"; |
| 8 | import { loadIdxView } from "@/worker/git/object-store"; |
| 9 | import { resolveDeltasAndWriteIdx, scanPack } from "@/worker/git/pack/indexer"; |
| 10 | import { loadPackRefView } from "@/worker/git/pack/refIndex"; |
| 11 | import { createQueueTaskContext, logSoftBudgetExhausted, retryQueueMessage } from "./context"; |
| 12 | |
| 13 | const REF_BACKFILL_SUBREQUEST_BUDGET = 7_500; |
| 14 | const REF_BACKFILL_RETRY_DELAY_SECONDS = 30; |
| 15 | |
| 16 | function 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 | |
| 26 | function 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 | |
| 37 | export 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 | } |