Skip to content
File

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

typescript424 lines
1import type { Logger } from "@/worker/common/logger";
2import type { PackCatalogRow } from "../../db/schema";
3import type { RepoLease, RepoStateSchema, TypedStorage } from "../../repoState";
4 
5import { asTypedStorage } from "../../repoState";
6import { scheduleAlarmIfSooner } from "../../scheduler";
7import { getActivePackCatalogSnapshot } from "../state";
8import { COMPACTION_WAKE_DELAY_MS, ensureRepoMetadataDefaults } from "../shared";
9 
10const COMPACTION_FAN_IN = 4;
11export const AUTO_COMPACTION_MAX_SOURCE_OBJECTS = 20_000;
12export const AUTO_COMPACTION_MAX_SOURCE_BYTES = 16 * 1024 * 1024;
13 
14// ---------------------------------------------------------------------------
15// Types
16// ---------------------------------------------------------------------------
17 
18export type CompactionTierState = {
19 tier: number;
20 activePackCount: number;
21};
22 
23export type CompactionPlan = {
24 sourceTier: number;
25 targetTier: number;
26 sourcePacks: PackCatalogRow[];
27 sourceBytes: number;
28 sourceObjects: number;
29 tiers: CompactionTierState[];
30};
31 
32export type CompactionBlockedReason = "over-budget" | "non-contiguous-window";
33 
34export type CompactionBlockedContext = {
35 reason: CompactionBlockedReason;
36 sourceTier: number;
37 activePackCount: number;
38 fanIn: number;
39 maxSourceObjects: number;
40 maxSourceBytes: number;
41 smallestWindowObjects?: number;
42 smallestWindowBytes?: number;
43 smallestWindowPackKeys?: string[];
44};
45 
46export type CompactionSelection =
47 | { status: "ready"; plan: CompactionPlan }
48 | { status: "no_work"; reason: "below-threshold"; tiers: CompactionTierState[] }
49 | { status: "blocked"; blocked: CompactionBlockedContext; tiers: CompactionTierState[] };
50 
51export type PreviewCompactionResult = {
52 action: "preview";
53 status: "ok" | "no_work" | "blocked";
54 queued: boolean;
55 wantedAt?: number;
56 activeCatalog: PackCatalogRow[];
57 packCatalogVersion: number;
58 plan?: CompactionPlan;
59 blocked?: CompactionBlockedContext;
60 reason?: "below-threshold" | CompactionBlockedReason;
61 message: string;
62};
63 
64export type RequestCompactionResult = {
65 action: "request";
66 status: "queued" | "no_work" | "blocked";
67 queued: boolean;
68 shouldEnqueue: boolean;
69 wantedAt?: number;
70 activeCatalog: PackCatalogRow[];
71 packCatalogVersion: number;
72 plan?: CompactionPlan;
73 blocked?: CompactionBlockedContext;
74 reason?: "below-threshold" | CompactionBlockedReason;
75 message: string;
76};
77 
78export type ClearCompactionRequestResult = {
79 action: "cleared";
80 cleared: boolean;
81 message: string;
82};
83 
84export type BeginCompactionResult =
85 | {
86 ok: true;
87 lease: RepoLease;
88 packsetVersion: number;
89 activeCatalog: PackCatalogRow[];
90 sourcePacks: PackCatalogRow[];
91 targetTier: number;
92 }
93 | {
94 ok: false;
95 status: "busy";
96 retryAfter: number;
97 reason: "receive-active" | "compact-active";
98 message: string;
99 }
100 | {
101 ok: false;
102 status: "no_work";
103 reason: "not-requested" | "below-threshold" | CompactionBlockedReason;
104 blocked?: CompactionBlockedContext;
105 message: string;
106 };
107 
108export type CommitCompactionResult =
109 | {
110 status: "committed";
111 packCatalogVersion: number;
112 shouldRequeue: boolean;
113 supersededPackKeys: string[];
114 targetPackKey: string;
115 }
116 | {
117 status: "retry";
118 reason: "receive-active" | "lease-mismatch" | "packset-changed" | "source-changed";
119 message: string;
120 };
121 
122// ---------------------------------------------------------------------------
123// Plan selection
124// ---------------------------------------------------------------------------
125 
126function summarizeTierCounts(activeCatalog: PackCatalogRow[]): CompactionTierState[] {
127 const counts = new Map<number, number>();
128 for (const pack of activeCatalog) {
129 counts.set(pack.tier, (counts.get(pack.tier) || 0) + 1);
130 }
131 
132 return Array.from(counts.entries())
133 .map(([tier, activePackCount]) => ({ tier, activePackCount }))
134 .sort((left, right) => left.tier - right.tier);
135}
136 
137function sumSourceBytes(sourcePacks: PackCatalogRow[]): number {
138 let total = 0;
139 for (const pack of sourcePacks) total += pack.packBytes;
140 return total;
141}
142 
143function sumSourceObjects(sourcePacks: PackCatalogRow[]): number {
144 let total = 0;
145 for (const pack of sourcePacks) total += pack.objectCount;
146 return total;
147}
148 
149type CompactionWindow = {
150 sourcePacks: PackCatalogRow[];
151 sourceBytes: number;
152 sourceObjects: number;
153 seqLo: number;
154 seqHi: number;
155};
156 
157type CompactionOverlapIndex = {
158 seqLoValues: number[];
159 seqHiValues: number[];
160};
161 
162function makeCompactionWindow(sourcePacks: PackCatalogRow[]): CompactionWindow {
163 let seqLo = sourcePacks[0]!.seqLo;
164 let seqHi = sourcePacks[0]!.seqHi;
165 for (const pack of sourcePacks) {
166 if (pack.seqLo < seqLo) seqLo = pack.seqLo;
167 if (pack.seqHi > seqHi) seqHi = pack.seqHi;
168 }
169 
170 return {
171 sourcePacks,
172 sourceBytes: sumSourceBytes(sourcePacks),
173 sourceObjects: sumSourceObjects(sourcePacks),
174 seqLo,
175 seqHi,
176 };
177}
178 
179function lowerBound(values: number[], target: number): number {
180 let lo = 0;
181 let hi = values.length;
182 while (lo < hi) {
183 const mid = Math.floor((lo + hi) / 2);
184 if (values[mid]! < target) {
185 lo = mid + 1;
186 } else {
187 hi = mid;
188 }
189 }
190 return lo;
191}
192 
193function upperBound(values: number[], target: number): number {
194 let lo = 0;
195 let hi = values.length;
196 while (lo < hi) {
197 const mid = Math.floor((lo + hi) / 2);
198 if (values[mid]! <= target) {
199 lo = mid + 1;
200 } else {
201 hi = mid;
202 }
203 }
204 return lo;
205}
206 
207function buildCompactionOverlapIndex(activeCatalog: PackCatalogRow[]): CompactionOverlapIndex {
208 return {
209 seqLoValues: activeCatalog.map((pack) => pack.seqLo).sort((left, right) => left - right),
210 seqHiValues: activeCatalog.map((pack) => pack.seqHi).sort((left, right) => left - right),
211 };
212}
213 
214function countOverlappingPacks(index: CompactionOverlapIndex, window: CompactionWindow): number {
215 const startedBeforeWindowEnd = upperBound(index.seqLoValues, window.seqHi);
216 const endedBeforeWindowStart = lowerBound(index.seqHiValues, window.seqLo);
217 return startedBeforeWindowEnd - endedBeforeWindowStart;
218}
219 
220function isClosedCompactionWindow(
221 overlapIndex: CompactionOverlapIndex,
222 window: CompactionWindow
223): boolean {
224 // Every selected source pack overlaps the window by construction. If the
225 // catalog has more overlaps than selected packs, an unselected active pack
226 // from this or another tier still covers part of the sequence range.
227 return countOverlappingPacks(overlapIndex, window) === window.sourcePacks.length;
228}
229 
230function withinAutomaticCompactionBudget(window: CompactionWindow): boolean {
231 return (
232 window.sourceObjects <= AUTO_COMPACTION_MAX_SOURCE_OBJECTS &&
233 window.sourceBytes <= AUTO_COMPACTION_MAX_SOURCE_BYTES
234 );
235}
236 
237function compareWindowCost(left: CompactionWindow, right: CompactionWindow): number {
238 const objectDiff = left.sourceObjects - right.sourceObjects;
239 if (objectDiff !== 0) return objectDiff;
240 return left.sourceBytes - right.sourceBytes;
241}
242 
243function buildPlanFromWindow(
244 window: CompactionWindow,
245 sourceTier: number,
246 tiers: CompactionTierState[]
247): CompactionPlan {
248 return {
249 sourceTier,
250 targetTier: sourceTier + 1,
251 sourcePacks: window.sourcePacks,
252 sourceBytes: window.sourceBytes,
253 sourceObjects: window.sourceObjects,
254 tiers,
255 };
256}
257 
258function buildBlockedContext(args: {
259 reason: CompactionBlockedReason;
260 sourceTier: number;
261 activePackCount: number;
262 smallestWindow?: CompactionWindow;
263}): CompactionBlockedContext {
264 return {
265 reason: args.reason,
266 sourceTier: args.sourceTier,
267 activePackCount: args.activePackCount,
268 fanIn: COMPACTION_FAN_IN,
269 maxSourceObjects: AUTO_COMPACTION_MAX_SOURCE_OBJECTS,
270 maxSourceBytes: AUTO_COMPACTION_MAX_SOURCE_BYTES,
271 smallestWindowObjects: args.smallestWindow?.sourceObjects,
272 smallestWindowBytes: args.smallestWindow?.sourceBytes,
273 smallestWindowPackKeys: args.smallestWindow?.sourcePacks.map((pack) => pack.packKey),
274 };
275}
276 
277export function selectCompactionWork(activeCatalog: PackCatalogRow[]): CompactionSelection {
278 const tiers = summarizeTierCounts(activeCatalog);
279 const overflowingTiers = tiers.filter((tier) => tier.activePackCount > COMPACTION_FAN_IN);
280 if (overflowingTiers.length === 0) {
281 return { status: "no_work", reason: "below-threshold", tiers };
282 }
283 
284 let blockedTier = overflowingTiers[0]!;
285 let smallestClosedWindow: CompactionWindow | undefined;
286 let sawNonClosedWindow = false;
287 const overlapIndex = buildCompactionOverlapIndex(activeCatalog);
288 
289 for (const overflowingTier of overflowingTiers) {
290 const tierPacks = activeCatalog
291 .filter((pack) => pack.tier === overflowingTier.tier)
292 .sort((left, right) => {
293 const seqLoDiff = left.seqLo - right.seqLo;
294 if (seqLoDiff !== 0) return seqLoDiff;
295 return left.seqHi - right.seqHi;
296 });
297 
298 for (let start = 0; start <= tierPacks.length - COMPACTION_FAN_IN; start++) {
299 const window = makeCompactionWindow(tierPacks.slice(start, start + COMPACTION_FAN_IN));
300 if (!isClosedCompactionWindow(overlapIndex, window)) {
301 sawNonClosedWindow = true;
302 continue;
303 }
304 
305 if (!smallestClosedWindow || compareWindowCost(window, smallestClosedWindow) < 0) {
306 smallestClosedWindow = window;
307 blockedTier = overflowingTier;
308 }
309 
310 // Automatic compaction is intentionally bounded. Large closed windows
311 // remain readable in-place and require a separate explicit full-repack
312 // path rather than letting queue maintenance exceed Worker CPU limits.
313 if (withinAutomaticCompactionBudget(window)) {
314 return {
315 status: "ready",
316 plan: buildPlanFromWindow(window, overflowingTier.tier, tiers),
317 };
318 }
319 }
320 }
321 
322 if (smallestClosedWindow) {
323 return {
324 status: "blocked",
325 blocked: buildBlockedContext({
326 reason: "over-budget",
327 sourceTier: blockedTier.tier,
328 activePackCount: blockedTier.activePackCount,
329 smallestWindow: smallestClosedWindow,
330 }),
331 tiers,
332 };
333 }
334 
335 return {
336 status: "blocked",
337 blocked: buildBlockedContext({
338 reason: sawNonClosedWindow ? "non-contiguous-window" : "over-budget",
339 sourceTier: blockedTier.tier,
340 activePackCount: blockedTier.activePackCount,
341 }),
342 tiers,
343 };
344}
345 
346export function selectCompactionPlan(activeCatalog: PackCatalogRow[]): CompactionPlan | undefined {
347 const selection = selectCompactionWork(activeCatalog);
348 return selection.status === "ready" ? selection.plan : undefined;
349}
350 
351export function catalogNeedsCompaction(activeCatalog: PackCatalogRow[]): boolean {
352 return selectCompactionWork(activeCatalog).status === "ready";
353}
354 
355// ---------------------------------------------------------------------------
356// Shared helpers used by requests.ts and lease.ts
357// ---------------------------------------------------------------------------
358 
359/** Returns true if every source pack in the commit request still matches the current catalog. */
360export function rowsMatchForCommit(
361 sourcePacks: PackCatalogRow[],
362 currentRows: PackCatalogRow[]
363): boolean {
364 if (sourcePacks.length !== currentRows.length) return false;
365 for (let index = 0; index < sourcePacks.length; index++) {
366 const expected = sourcePacks[index];
367 const current = currentRows[index];
368 if (!current) return false;
369 if (current.packKey !== expected.packKey) return false;
370 if (current.state !== "active") return false;
371 if (current.kind !== expected.kind) return false;
372 if (current.tier !== expected.tier) return false;
373 if (current.seqLo !== expected.seqLo || current.seqHi !== expected.seqHi) return false;
374 if (current.objectCount !== expected.objectCount) return false;
375 if (current.packBytes !== expected.packBytes || current.idxBytes !== expected.idxBytes)
376 return false;
377 }
378 return true;
379}
380 
381export async function scheduleCompactionAlarm(
382 ctx: DurableObjectState,
383 env: Env,
384 delayMs: number
385): Promise<void> {
386 await scheduleAlarmIfSooner(ctx, env, Date.now() + delayMs);
387}
388 
389export async function scheduleCompactionWake(ctx: DurableObjectState, env: Env): Promise<void> {
390 await scheduleCompactionAlarm(ctx, env, COMPACTION_WAKE_DELAY_MS);
391}
392 
393/** Load shared compaction context: catalog, plan, and queued state. */
394export async function loadCompactionContext(args: {
395 ctx: DurableObjectState;
396 env: Env;
397 prefix: string;
398 logger?: Logger;
399}): Promise<{
400 store: TypedStorage<RepoStateSchema>;
401 packCatalogVersion: number;
402 wantedAt: number | undefined;
403 activeCatalog: PackCatalogRow[];
404 plan: CompactionPlan | undefined;
405 selection: CompactionSelection;
406}> {
407 const store = asTypedStorage<RepoStateSchema>(args.ctx.storage);
408 await ensureRepoMetadataDefaults(store);
409 const packCatalogVersion = (await store.get("packsetVersion")) || 0;
410 const wantedAt = await store.get("compactionWantedAt");
411 const activeCatalog = await getActivePackCatalogSnapshot(args.ctx);
412 const selection = selectCompactionWork(activeCatalog);
413 const plan = selection.status === "ready" ? selection.plan : undefined;
414 
415 return {
416 store,
417 packCatalogVersion,
418 wantedAt,
419 activeCatalog,
420 plan,
421 selection,
422 };
423}