File
Blob: src/worker/tasks/compaction.ts
| 1 | import type { CacheContext } from "@/worker/cache"; |
| 2 | import type { Logger } from "@/worker/common/logger"; |
| 3 | import type { OrderedPackSnapshot } from "@/worker/git/operations/fetch/types"; |
| 4 | import type { RepoDurableObject } from "@/worker/do/repo/repoDO"; |
| 5 | |
| 6 | import { |
| 7 | type CompactionDeleteQueueMessage, |
| 8 | type CompactionQueueMessage, |
| 9 | type RepoQueueMessageHandle, |
| 10 | } from "./types"; |
| 11 | |
| 12 | import { getRepoStubByDoId } from "@/worker/common"; |
| 13 | import { buildCompactionNeededOids } from "@/worker/git/compaction/plan"; |
| 14 | import { type SubrequestLimiter } from "@/worker/git/operations/limits"; |
| 15 | import { scanPack, resolveDeltasAndWriteIdx } from "@/worker/git/pack/indexer"; |
| 16 | import { rewritePackResult } from "@/worker/git/pack/rewrite"; |
| 17 | import { loadOrderedPackSnapshot } from "@/worker/git/pack/snapshot"; |
| 18 | import { |
| 19 | deleteStagedPack, |
| 20 | stagePackToR2, |
| 21 | type StagedPackUpload, |
| 22 | } from "@/worker/git/receive/r2Upload"; |
| 23 | import { doPrefix, packIndexKey, packRefsKey, r2PackKey } from "@/worker/keys"; |
| 24 | import { createQueueTaskContext, logSoftBudgetExhausted, retryQueueMessage } from "./context"; |
| 25 | |
| 26 | const COMPACTION_SUBREQUEST_BUDGET = 7_500; |
| 27 | const COMPACTION_RETRY_DELAY_SECONDS = 30; |
| 28 | const COMPACTION_CONFLICT_RETRY_DELAY_SECONDS = 10; |
| 29 | const COMPACTION_DELETE_DELAY_SECONDS = 60; |
| 30 | |
| 31 | function 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 | |
| 41 | async 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 | |
| 58 | async 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 | |
| 93 | async 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 | |
| 114 | export 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 | |
| 377 | export 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 | } |