File
Blob: src/worker/do/repo/catalog/compaction/lease.ts
| 1 | /** |
| 2 | * Queue-facing compaction state transitions: begin, commit, abort, and alarm rearm. |
| 3 | * |
| 4 | * These operations manage the compaction lease lifecycle. The queue consumer |
| 5 | * acquires a lease via `beginCompactionState`, performs the pack rewrite in |
| 6 | * worker code, and then atomically commits the result via `commitCompactionState`. |
| 7 | */ |
| 8 | import type { Logger } from "@/worker/common/logger"; |
| 9 | import type { PackCatalogRow } from "../../db/schema"; |
| 10 | import type { RepoLease, RepoStateSchema } from "../../repoState"; |
| 11 | |
| 12 | import { asTypedStorage } from "../../repoState"; |
| 13 | import { |
| 14 | getDb, |
| 15 | getPackCatalogRow, |
| 16 | listActivePackCatalog, |
| 17 | supersedePackCatalogRows, |
| 18 | upsertPackCatalogRow, |
| 19 | } from "../../db"; |
| 20 | import { clearExpiredLeases } from "../leases"; |
| 21 | import { getActivePackCatalogSnapshot } from "../state"; |
| 22 | import { |
| 23 | bumpPacksetVersion, |
| 24 | COMPACT_LEASE_TTL_MS, |
| 25 | COMPACTION_REARM_DELAY_MS, |
| 26 | ensureRepoMetadataDefaults, |
| 27 | LEASE_RETRY_AFTER_SECONDS, |
| 28 | } from "../shared"; |
| 29 | import { activeLeaseOrUndefined } from "../activity"; |
| 30 | import { |
| 31 | selectCompactionWork, |
| 32 | scheduleCompactionWake, |
| 33 | scheduleCompactionAlarm, |
| 34 | rowsMatchForCommit, |
| 35 | type BeginCompactionResult, |
| 36 | type CommitCompactionResult, |
| 37 | } from "./plan"; |
| 38 | |
| 39 | /** |
| 40 | * Acquire a compaction lease and select source packs for compaction. |
| 41 | * |
| 42 | * Rejects when: no compaction request is recorded, a receive or compaction |
| 43 | * lease is already active, the catalog is already within policy, or every |
| 44 | * automatic source window is too expensive for bounded queue maintenance. |
| 45 | */ |
| 46 | export async function beginCompactionState(args: { |
| 47 | ctx: DurableObjectState; |
| 48 | env: Env; |
| 49 | prefix: string; |
| 50 | logger?: Logger; |
| 51 | }): Promise<BeginCompactionResult> { |
| 52 | const store = asTypedStorage<RepoStateSchema>(args.ctx.storage); |
| 53 | await clearExpiredLeases(args.ctx, args.logger); |
| 54 | await ensureRepoMetadataDefaults(store); |
| 55 | |
| 56 | const wantedAt = await store.get("compactionWantedAt"); |
| 57 | if (typeof wantedAt !== "number") { |
| 58 | return { |
| 59 | ok: false, |
| 60 | status: "no_work", |
| 61 | reason: "not-requested", |
| 62 | message: "No compaction request is currently recorded for this repository.", |
| 63 | }; |
| 64 | } |
| 65 | |
| 66 | const now = Date.now(); |
| 67 | const receiveLease = activeLeaseOrUndefined(await store.get("receiveLease"), now); |
| 68 | if (receiveLease) { |
| 69 | return { |
| 70 | ok: false, |
| 71 | status: "busy", |
| 72 | retryAfter: LEASE_RETRY_AFTER_SECONDS, |
| 73 | reason: "receive-active", |
| 74 | message: "A receive lease is active, so compaction must retry later.", |
| 75 | }; |
| 76 | } |
| 77 | |
| 78 | const compactLease = activeLeaseOrUndefined(await store.get("compactLease"), now); |
| 79 | if (compactLease) { |
| 80 | return { |
| 81 | ok: false, |
| 82 | status: "busy", |
| 83 | retryAfter: LEASE_RETRY_AFTER_SECONDS, |
| 84 | reason: "compact-active", |
| 85 | message: "A compaction lease is already active for this repository.", |
| 86 | }; |
| 87 | } |
| 88 | |
| 89 | const activeCatalog = await getActivePackCatalogSnapshot(args.ctx); |
| 90 | const selection = selectCompactionWork(activeCatalog); |
| 91 | if (selection.status === "blocked") { |
| 92 | await store.delete("compactionWantedAt"); |
| 93 | args.logger?.warn("compaction:begin-blocked", { |
| 94 | reason: selection.blocked.reason, |
| 95 | sourceTier: selection.blocked.sourceTier, |
| 96 | activePackCount: selection.blocked.activePackCount, |
| 97 | maxSourceObjects: selection.blocked.maxSourceObjects, |
| 98 | maxSourceBytes: selection.blocked.maxSourceBytes, |
| 99 | smallestWindowObjects: selection.blocked.smallestWindowObjects, |
| 100 | smallestWindowBytes: selection.blocked.smallestWindowBytes, |
| 101 | }); |
| 102 | return { |
| 103 | ok: false, |
| 104 | status: "no_work", |
| 105 | reason: selection.blocked.reason, |
| 106 | blocked: selection.blocked, |
| 107 | message: |
| 108 | "Automatic compaction is blocked because no safe source-pack window fits the bounded maintenance budget.", |
| 109 | }; |
| 110 | } |
| 111 | |
| 112 | if (selection.status === "no_work") { |
| 113 | await store.delete("compactionWantedAt"); |
| 114 | args.logger?.info("compaction:begin-no-work", { |
| 115 | reason: "below-threshold", |
| 116 | }); |
| 117 | return { |
| 118 | ok: false, |
| 119 | status: "no_work", |
| 120 | reason: "below-threshold", |
| 121 | message: "The active pack catalog is already within the compaction policy.", |
| 122 | }; |
| 123 | } |
| 124 | const plan = selection.plan; |
| 125 | |
| 126 | const lease: RepoLease = { |
| 127 | token: crypto.randomUUID(), |
| 128 | createdAt: now, |
| 129 | expiresAt: now + COMPACT_LEASE_TTL_MS, |
| 130 | }; |
| 131 | await store.put("compactLease", lease); |
| 132 | |
| 133 | args.logger?.info("compaction:begin", { |
| 134 | leaseToken: lease.token, |
| 135 | sourceTier: plan.sourceTier, |
| 136 | targetTier: plan.targetTier, |
| 137 | sourceCount: plan.sourcePacks.length, |
| 138 | }); |
| 139 | return { |
| 140 | ok: true, |
| 141 | lease, |
| 142 | packsetVersion: (await store.get("packsetVersion")) || 0, |
| 143 | activeCatalog, |
| 144 | sourcePacks: plan.sourcePacks, |
| 145 | targetTier: plan.targetTier, |
| 146 | }; |
| 147 | } |
| 148 | |
| 149 | /** |
| 150 | * Atomically commit a compaction result: insert the new pack, supersede source |
| 151 | * packs, bump the packset version, and mirror legacy keys. |
| 152 | * |
| 153 | * Rejects with `status: "retry"` when the lease is stale, a receive lease |
| 154 | * appeared, the packset version changed, or source packs were modified since |
| 155 | * `beginCompactionState`. |
| 156 | */ |
| 157 | export async function commitCompactionState(args: { |
| 158 | ctx: DurableObjectState; |
| 159 | env: Env; |
| 160 | token: string; |
| 161 | sourcePacks: PackCatalogRow[]; |
| 162 | targetTier: number; |
| 163 | packsetVersion: number; |
| 164 | stagedPack: { |
| 165 | packKey: string; |
| 166 | packBytes: number; |
| 167 | idxBytes: number; |
| 168 | objectCount: number; |
| 169 | }; |
| 170 | logger?: Logger; |
| 171 | }): Promise<CommitCompactionResult> { |
| 172 | const store = asTypedStorage<RepoStateSchema>(args.ctx.storage); |
| 173 | await ensureRepoMetadataDefaults(store); |
| 174 | |
| 175 | const lease = await store.get("compactLease"); |
| 176 | if (!lease || lease.token !== args.token) { |
| 177 | return { |
| 178 | status: "retry", |
| 179 | reason: "lease-mismatch", |
| 180 | message: "Compaction lease is no longer active for this request.", |
| 181 | }; |
| 182 | } |
| 183 | |
| 184 | const receiveLease = activeLeaseOrUndefined(await store.get("receiveLease"), Date.now()); |
| 185 | if (receiveLease) { |
| 186 | await store.delete("compactLease"); |
| 187 | return { |
| 188 | status: "retry", |
| 189 | reason: "receive-active", |
| 190 | message: "A receive lease became active before compaction could commit.", |
| 191 | }; |
| 192 | } |
| 193 | |
| 194 | const currentPacksetVersion = (await store.get("packsetVersion")) || 0; |
| 195 | if (currentPacksetVersion !== args.packsetVersion) { |
| 196 | await store.delete("compactLease"); |
| 197 | return { |
| 198 | status: "retry", |
| 199 | reason: "packset-changed", |
| 200 | message: "The active pack catalog changed before compaction could commit.", |
| 201 | }; |
| 202 | } |
| 203 | |
| 204 | const db = getDb(args.ctx.storage); |
| 205 | const currentRows: PackCatalogRow[] = []; |
| 206 | for (const sourcePack of args.sourcePacks) { |
| 207 | const row = await getPackCatalogRow(db, sourcePack.packKey); |
| 208 | if (row) currentRows.push(row); |
| 209 | } |
| 210 | if (!rowsMatchForCommit(args.sourcePacks, currentRows)) { |
| 211 | await store.delete("compactLease"); |
| 212 | return { |
| 213 | status: "retry", |
| 214 | reason: "source-changed", |
| 215 | message: "One or more source packs changed before compaction could commit.", |
| 216 | }; |
| 217 | } |
| 218 | |
| 219 | let seqLo = args.sourcePacks[0]!.seqLo; |
| 220 | let seqHi = args.sourcePacks[0]!.seqHi; |
| 221 | for (const sourcePack of args.sourcePacks) { |
| 222 | if (sourcePack.seqLo < seqLo) seqLo = sourcePack.seqLo; |
| 223 | if (sourcePack.seqHi > seqHi) seqHi = sourcePack.seqHi; |
| 224 | } |
| 225 | |
| 226 | await upsertPackCatalogRow(db, { |
| 227 | packKey: args.stagedPack.packKey, |
| 228 | kind: "compact", |
| 229 | state: "active", |
| 230 | tier: args.targetTier, |
| 231 | seqLo, |
| 232 | seqHi, |
| 233 | objectCount: args.stagedPack.objectCount, |
| 234 | packBytes: args.stagedPack.packBytes, |
| 235 | idxBytes: args.stagedPack.idxBytes, |
| 236 | createdAt: Date.now(), |
| 237 | supersededBy: null, |
| 238 | }); |
| 239 | await supersedePackCatalogRows( |
| 240 | db, |
| 241 | args.sourcePacks.map((row) => row.packKey), |
| 242 | args.stagedPack.packKey |
| 243 | ); |
| 244 | |
| 245 | const activeCatalog = await listActivePackCatalog(db); |
| 246 | const nextPackCatalogVersion = await bumpPacksetVersion(store); |
| 247 | |
| 248 | const compactionSelection = selectCompactionWork(activeCatalog); |
| 249 | const shouldRequeue = compactionSelection.status === "ready"; |
| 250 | if (shouldRequeue) { |
| 251 | await store.put("compactionWantedAt", Date.now()); |
| 252 | await scheduleCompactionWake(args.ctx, args.env); |
| 253 | } else { |
| 254 | await store.delete("compactionWantedAt"); |
| 255 | if (compactionSelection.status === "blocked") { |
| 256 | args.logger?.warn("compaction:commit-requeue-blocked", { |
| 257 | reason: compactionSelection.blocked.reason, |
| 258 | sourceTier: compactionSelection.blocked.sourceTier, |
| 259 | activePackCount: compactionSelection.blocked.activePackCount, |
| 260 | maxSourceObjects: compactionSelection.blocked.maxSourceObjects, |
| 261 | maxSourceBytes: compactionSelection.blocked.maxSourceBytes, |
| 262 | smallestWindowObjects: compactionSelection.blocked.smallestWindowObjects, |
| 263 | smallestWindowBytes: compactionSelection.blocked.smallestWindowBytes, |
| 264 | }); |
| 265 | } |
| 266 | } |
| 267 | |
| 268 | await store.delete("compactLease"); |
| 269 | args.logger?.info("compaction:commit", { |
| 270 | targetPackKey: args.stagedPack.packKey, |
| 271 | supersededCount: args.sourcePacks.length, |
| 272 | shouldRequeue, |
| 273 | packCatalogVersion: nextPackCatalogVersion, |
| 274 | }); |
| 275 | return { |
| 276 | status: "committed", |
| 277 | packCatalogVersion: nextPackCatalogVersion, |
| 278 | shouldRequeue, |
| 279 | supersededPackKeys: args.sourcePacks.map((row) => row.packKey), |
| 280 | targetPackKey: args.stagedPack.packKey, |
| 281 | }; |
| 282 | } |
| 283 | |
| 284 | /** |
| 285 | * Called from the DO alarm handler for streaming repos. If `compactionWantedAt` |
| 286 | * is set and no leases are active, enqueue a compaction message to the |
| 287 | * maintenance queue. Reschedules the alarm on queue send failure. |
| 288 | */ |
| 289 | export async function rearmCompactionQueueFromAlarm(args: { |
| 290 | ctx: DurableObjectState; |
| 291 | env: Env; |
| 292 | logger?: Logger; |
| 293 | }): Promise<boolean> { |
| 294 | const store = asTypedStorage<RepoStateSchema>(args.ctx.storage); |
| 295 | await ensureRepoMetadataDefaults(store); |
| 296 | |
| 297 | const wantedAt = await store.get("compactionWantedAt"); |
| 298 | if (typeof wantedAt !== "number") return false; |
| 299 | |
| 300 | const now = Date.now(); |
| 301 | if (activeLeaseOrUndefined(await store.get("receiveLease"), now)) return false; |
| 302 | if (activeLeaseOrUndefined(await store.get("compactLease"), now)) return false; |
| 303 | |
| 304 | const doId = args.ctx.id.toString(); |
| 305 | try { |
| 306 | await args.env.REPO_TASKS_QUEUE.send({ |
| 307 | kind: "compaction", |
| 308 | doId, |
| 309 | }); |
| 310 | args.logger?.info("compaction:alarm-rearm-enqueued", { doId }); |
| 311 | return true; |
| 312 | } catch (error) { |
| 313 | args.logger?.warn("compaction:alarm-rearm-failed", { |
| 314 | doId, |
| 315 | error: String(error), |
| 316 | }); |
| 317 | await scheduleCompactionAlarm(args.ctx, args.env, COMPACTION_REARM_DELAY_MS); |
| 318 | return true; |
| 319 | } |
| 320 | } |