Skip to content
File

Blob: src/worker/do/repo/catalog/compaction/lease.ts

typescript321 lines
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 */
8import type { Logger } from "@/worker/common/logger";
9import type { PackCatalogRow } from "../../db/schema";
10import type { RepoLease, RepoStateSchema } from "../../repoState";
11 
12import { asTypedStorage } from "../../repoState";
13import {
14 getDb,
15 getPackCatalogRow,
16 listActivePackCatalog,
17 supersedePackCatalogRows,
18 upsertPackCatalogRow,
19} from "../../db";
20import { clearExpiredLeases } from "../leases";
21import { getActivePackCatalogSnapshot } from "../state";
22import {
23 bumpPacksetVersion,
24 COMPACT_LEASE_TTL_MS,
25 COMPACTION_REARM_DELAY_MS,
26 ensureRepoMetadataDefaults,
27 LEASE_RETRY_AFTER_SECONDS,
28} from "../shared";
29import { activeLeaseOrUndefined } from "../activity";
30import {
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 */
46export 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 */
157export 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 */
289export 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}