import type { Logger } from "@/worker/common/logger"; import type { PackCatalogRow } from "../../db/schema"; import type { RepoLease, RepoStateSchema, TypedStorage } from "../../repoState"; import { asTypedStorage } from "../../repoState"; import { scheduleAlarmIfSooner } from "../../scheduler"; import { getActivePackCatalogSnapshot } from "../state"; import { COMPACTION_WAKE_DELAY_MS, ensureRepoMetadataDefaults } from "../shared"; const COMPACTION_FAN_IN = 4; export const AUTO_COMPACTION_MAX_SOURCE_OBJECTS = 20_000; export const AUTO_COMPACTION_MAX_SOURCE_BYTES = 16 * 1024 * 1024; // --------------------------------------------------------------------------- // Types // --------------------------------------------------------------------------- export type CompactionTierState = { tier: number; activePackCount: number; }; export type CompactionPlan = { sourceTier: number; targetTier: number; sourcePacks: PackCatalogRow[]; sourceBytes: number; sourceObjects: number; tiers: CompactionTierState[]; }; export type CompactionBlockedReason = "over-budget" | "non-contiguous-window"; export type CompactionBlockedContext = { reason: CompactionBlockedReason; sourceTier: number; activePackCount: number; fanIn: number; maxSourceObjects: number; maxSourceBytes: number; smallestWindowObjects?: number; smallestWindowBytes?: number; smallestWindowPackKeys?: string[]; }; export type CompactionSelection = | { status: "ready"; plan: CompactionPlan } | { status: "no_work"; reason: "below-threshold"; tiers: CompactionTierState[] } | { status: "blocked"; blocked: CompactionBlockedContext; tiers: CompactionTierState[] }; export type PreviewCompactionResult = { action: "preview"; status: "ok" | "no_work" | "blocked"; queued: boolean; wantedAt?: number; activeCatalog: PackCatalogRow[]; packCatalogVersion: number; plan?: CompactionPlan; blocked?: CompactionBlockedContext; reason?: "below-threshold" | CompactionBlockedReason; message: string; }; export type RequestCompactionResult = { action: "request"; status: "queued" | "no_work" | "blocked"; queued: boolean; shouldEnqueue: boolean; wantedAt?: number; activeCatalog: PackCatalogRow[]; packCatalogVersion: number; plan?: CompactionPlan; blocked?: CompactionBlockedContext; reason?: "below-threshold" | CompactionBlockedReason; message: string; }; export type ClearCompactionRequestResult = { action: "cleared"; cleared: boolean; message: string; }; export type BeginCompactionResult = | { ok: true; lease: RepoLease; packsetVersion: number; activeCatalog: PackCatalogRow[]; sourcePacks: PackCatalogRow[]; targetTier: number; } | { ok: false; status: "busy"; retryAfter: number; reason: "receive-active" | "compact-active"; message: string; } | { ok: false; status: "no_work"; reason: "not-requested" | "below-threshold" | CompactionBlockedReason; blocked?: CompactionBlockedContext; message: string; }; export type CommitCompactionResult = | { status: "committed"; packCatalogVersion: number; shouldRequeue: boolean; supersededPackKeys: string[]; targetPackKey: string; } | { status: "retry"; reason: "receive-active" | "lease-mismatch" | "packset-changed" | "source-changed"; message: string; }; // --------------------------------------------------------------------------- // Plan selection // --------------------------------------------------------------------------- function summarizeTierCounts(activeCatalog: PackCatalogRow[]): CompactionTierState[] { const counts = new Map(); for (const pack of activeCatalog) { counts.set(pack.tier, (counts.get(pack.tier) || 0) + 1); } return Array.from(counts.entries()) .map(([tier, activePackCount]) => ({ tier, activePackCount })) .sort((left, right) => left.tier - right.tier); } function sumSourceBytes(sourcePacks: PackCatalogRow[]): number { let total = 0; for (const pack of sourcePacks) total += pack.packBytes; return total; } function sumSourceObjects(sourcePacks: PackCatalogRow[]): number { let total = 0; for (const pack of sourcePacks) total += pack.objectCount; return total; } type CompactionWindow = { sourcePacks: PackCatalogRow[]; sourceBytes: number; sourceObjects: number; seqLo: number; seqHi: number; }; type CompactionOverlapIndex = { seqLoValues: number[]; seqHiValues: number[]; }; function makeCompactionWindow(sourcePacks: PackCatalogRow[]): CompactionWindow { let seqLo = sourcePacks[0]!.seqLo; let seqHi = sourcePacks[0]!.seqHi; for (const pack of sourcePacks) { if (pack.seqLo < seqLo) seqLo = pack.seqLo; if (pack.seqHi > seqHi) seqHi = pack.seqHi; } return { sourcePacks, sourceBytes: sumSourceBytes(sourcePacks), sourceObjects: sumSourceObjects(sourcePacks), seqLo, seqHi, }; } function lowerBound(values: number[], target: number): number { let lo = 0; let hi = values.length; while (lo < hi) { const mid = Math.floor((lo + hi) / 2); if (values[mid]! < target) { lo = mid + 1; } else { hi = mid; } } return lo; } function upperBound(values: number[], target: number): number { let lo = 0; let hi = values.length; while (lo < hi) { const mid = Math.floor((lo + hi) / 2); if (values[mid]! <= target) { lo = mid + 1; } else { hi = mid; } } return lo; } function buildCompactionOverlapIndex(activeCatalog: PackCatalogRow[]): CompactionOverlapIndex { return { seqLoValues: activeCatalog.map((pack) => pack.seqLo).sort((left, right) => left - right), seqHiValues: activeCatalog.map((pack) => pack.seqHi).sort((left, right) => left - right), }; } function countOverlappingPacks(index: CompactionOverlapIndex, window: CompactionWindow): number { const startedBeforeWindowEnd = upperBound(index.seqLoValues, window.seqHi); const endedBeforeWindowStart = lowerBound(index.seqHiValues, window.seqLo); return startedBeforeWindowEnd - endedBeforeWindowStart; } function isClosedCompactionWindow( overlapIndex: CompactionOverlapIndex, window: CompactionWindow ): boolean { // Every selected source pack overlaps the window by construction. If the // catalog has more overlaps than selected packs, an unselected active pack // from this or another tier still covers part of the sequence range. return countOverlappingPacks(overlapIndex, window) === window.sourcePacks.length; } function withinAutomaticCompactionBudget(window: CompactionWindow): boolean { return ( window.sourceObjects <= AUTO_COMPACTION_MAX_SOURCE_OBJECTS && window.sourceBytes <= AUTO_COMPACTION_MAX_SOURCE_BYTES ); } function compareWindowCost(left: CompactionWindow, right: CompactionWindow): number { const objectDiff = left.sourceObjects - right.sourceObjects; if (objectDiff !== 0) return objectDiff; return left.sourceBytes - right.sourceBytes; } function buildPlanFromWindow( window: CompactionWindow, sourceTier: number, tiers: CompactionTierState[] ): CompactionPlan { return { sourceTier, targetTier: sourceTier + 1, sourcePacks: window.sourcePacks, sourceBytes: window.sourceBytes, sourceObjects: window.sourceObjects, tiers, }; } function buildBlockedContext(args: { reason: CompactionBlockedReason; sourceTier: number; activePackCount: number; smallestWindow?: CompactionWindow; }): CompactionBlockedContext { return { reason: args.reason, sourceTier: args.sourceTier, activePackCount: args.activePackCount, fanIn: COMPACTION_FAN_IN, maxSourceObjects: AUTO_COMPACTION_MAX_SOURCE_OBJECTS, maxSourceBytes: AUTO_COMPACTION_MAX_SOURCE_BYTES, smallestWindowObjects: args.smallestWindow?.sourceObjects, smallestWindowBytes: args.smallestWindow?.sourceBytes, smallestWindowPackKeys: args.smallestWindow?.sourcePacks.map((pack) => pack.packKey), }; } export function selectCompactionWork(activeCatalog: PackCatalogRow[]): CompactionSelection { const tiers = summarizeTierCounts(activeCatalog); const overflowingTiers = tiers.filter((tier) => tier.activePackCount > COMPACTION_FAN_IN); if (overflowingTiers.length === 0) { return { status: "no_work", reason: "below-threshold", tiers }; } let blockedTier = overflowingTiers[0]!; let smallestClosedWindow: CompactionWindow | undefined; let sawNonClosedWindow = false; const overlapIndex = buildCompactionOverlapIndex(activeCatalog); for (const overflowingTier of overflowingTiers) { const tierPacks = activeCatalog .filter((pack) => pack.tier === overflowingTier.tier) .sort((left, right) => { const seqLoDiff = left.seqLo - right.seqLo; if (seqLoDiff !== 0) return seqLoDiff; return left.seqHi - right.seqHi; }); for (let start = 0; start <= tierPacks.length - COMPACTION_FAN_IN; start++) { const window = makeCompactionWindow(tierPacks.slice(start, start + COMPACTION_FAN_IN)); if (!isClosedCompactionWindow(overlapIndex, window)) { sawNonClosedWindow = true; continue; } if (!smallestClosedWindow || compareWindowCost(window, smallestClosedWindow) < 0) { smallestClosedWindow = window; blockedTier = overflowingTier; } // Automatic compaction is intentionally bounded. Large closed windows // remain readable in-place and require a separate explicit full-repack // path rather than letting queue maintenance exceed Worker CPU limits. if (withinAutomaticCompactionBudget(window)) { return { status: "ready", plan: buildPlanFromWindow(window, overflowingTier.tier, tiers), }; } } } if (smallestClosedWindow) { return { status: "blocked", blocked: buildBlockedContext({ reason: "over-budget", sourceTier: blockedTier.tier, activePackCount: blockedTier.activePackCount, smallestWindow: smallestClosedWindow, }), tiers, }; } return { status: "blocked", blocked: buildBlockedContext({ reason: sawNonClosedWindow ? "non-contiguous-window" : "over-budget", sourceTier: blockedTier.tier, activePackCount: blockedTier.activePackCount, }), tiers, }; } export function selectCompactionPlan(activeCatalog: PackCatalogRow[]): CompactionPlan | undefined { const selection = selectCompactionWork(activeCatalog); return selection.status === "ready" ? selection.plan : undefined; } export function catalogNeedsCompaction(activeCatalog: PackCatalogRow[]): boolean { return selectCompactionWork(activeCatalog).status === "ready"; } // --------------------------------------------------------------------------- // Shared helpers used by requests.ts and lease.ts // --------------------------------------------------------------------------- /** Returns true if every source pack in the commit request still matches the current catalog. */ export function rowsMatchForCommit( sourcePacks: PackCatalogRow[], currentRows: PackCatalogRow[] ): boolean { if (sourcePacks.length !== currentRows.length) return false; for (let index = 0; index < sourcePacks.length; index++) { const expected = sourcePacks[index]; const current = currentRows[index]; if (!current) return false; if (current.packKey !== expected.packKey) return false; if (current.state !== "active") return false; if (current.kind !== expected.kind) return false; if (current.tier !== expected.tier) return false; if (current.seqLo !== expected.seqLo || current.seqHi !== expected.seqHi) return false; if (current.objectCount !== expected.objectCount) return false; if (current.packBytes !== expected.packBytes || current.idxBytes !== expected.idxBytes) return false; } return true; } export async function scheduleCompactionAlarm( ctx: DurableObjectState, env: Env, delayMs: number ): Promise { await scheduleAlarmIfSooner(ctx, env, Date.now() + delayMs); } export async function scheduleCompactionWake(ctx: DurableObjectState, env: Env): Promise { await scheduleCompactionAlarm(ctx, env, COMPACTION_WAKE_DELAY_MS); } /** Load shared compaction context: catalog, plan, and queued state. */ export async function loadCompactionContext(args: { ctx: DurableObjectState; env: Env; prefix: string; logger?: Logger; }): Promise<{ store: TypedStorage; packCatalogVersion: number; wantedAt: number | undefined; activeCatalog: PackCatalogRow[]; plan: CompactionPlan | undefined; selection: CompactionSelection; }> { const store = asTypedStorage(args.ctx.storage); await ensureRepoMetadataDefaults(store); const packCatalogVersion = (await store.get("packsetVersion")) || 0; const wantedAt = await store.get("compactionWantedAt"); const activeCatalog = await getActivePackCatalogSnapshot(args.ctx); const selection = selectCompactionWork(activeCatalog); const plan = selection.status === "ready" ? selection.plan : undefined; return { store, packCatalogVersion, wantedAt, activeCatalog, plan, selection, }; }