File
Blob: src/worker/do/repo/catalog/compaction/requests.ts
| 1 | /** |
| 2 | * Admin-facing compaction state transitions: preview, request, and clear. |
| 3 | * |
| 4 | * These operations are triggered by POST/DELETE on /admin/compact. |
| 5 | * They read or mutate `compactionWantedAt` in DO storage but never |
| 6 | * acquire or release compaction leases. |
| 7 | */ |
| 8 | import type { Logger } from "@/worker/common/logger"; |
| 9 | import type { RepoStateSchema } from "../../repoState"; |
| 10 | |
| 11 | import { asTypedStorage } from "../../repoState"; |
| 12 | import { |
| 13 | loadCompactionContext, |
| 14 | scheduleCompactionWake, |
| 15 | type PreviewCompactionResult, |
| 16 | type RequestCompactionResult, |
| 17 | type ClearCompactionRequestResult, |
| 18 | } from "./plan"; |
| 19 | |
| 20 | /** |
| 21 | * Preview the current compaction plan without recording a request. |
| 22 | * |
| 23 | * Returns a concrete automatic plan only when a compactable tier has a safe |
| 24 | * bounded source window. An overflowing tier can still report `blocked` when |
| 25 | * automatic queue maintenance would exceed the configured in-code budget. |
| 26 | */ |
| 27 | export async function previewCompactionState(args: { |
| 28 | ctx: DurableObjectState; |
| 29 | env: Env; |
| 30 | prefix: string; |
| 31 | logger?: Logger; |
| 32 | }): Promise<PreviewCompactionResult> { |
| 33 | const context = await loadCompactionContext(args); |
| 34 | const queued = typeof context.wantedAt === "number"; |
| 35 | |
| 36 | if (context.selection.status === "blocked") { |
| 37 | args.logger?.warn("compaction:preview-blocked", { |
| 38 | queued, |
| 39 | reason: context.selection.blocked.reason, |
| 40 | sourceTier: context.selection.blocked.sourceTier, |
| 41 | activePackCount: context.selection.blocked.activePackCount, |
| 42 | maxSourceObjects: context.selection.blocked.maxSourceObjects, |
| 43 | maxSourceBytes: context.selection.blocked.maxSourceBytes, |
| 44 | smallestWindowObjects: context.selection.blocked.smallestWindowObjects, |
| 45 | smallestWindowBytes: context.selection.blocked.smallestWindowBytes, |
| 46 | }); |
| 47 | return { |
| 48 | action: "preview", |
| 49 | status: "blocked", |
| 50 | queued, |
| 51 | wantedAt: context.wantedAt, |
| 52 | activeCatalog: context.activeCatalog, |
| 53 | packCatalogVersion: context.packCatalogVersion, |
| 54 | blocked: context.selection.blocked, |
| 55 | reason: context.selection.blocked.reason, |
| 56 | message: |
| 57 | "Automatic compaction is blocked because no safe source-pack window fits the bounded maintenance budget.", |
| 58 | }; |
| 59 | } |
| 60 | |
| 61 | if (context.selection.status === "no_work") { |
| 62 | args.logger?.info("compaction:preview-no-work", { |
| 63 | reason: "below-threshold", |
| 64 | queued, |
| 65 | }); |
| 66 | return { |
| 67 | action: "preview", |
| 68 | status: "no_work", |
| 69 | queued, |
| 70 | wantedAt: context.wantedAt, |
| 71 | activeCatalog: context.activeCatalog, |
| 72 | packCatalogVersion: context.packCatalogVersion, |
| 73 | reason: "below-threshold", |
| 74 | message: "The active pack catalog is already within the compaction policy.", |
| 75 | }; |
| 76 | } |
| 77 | |
| 78 | const plan = context.selection.plan; |
| 79 | args.logger?.info("compaction:preview", { |
| 80 | queued, |
| 81 | sourceTier: plan.sourceTier, |
| 82 | targetTier: plan.targetTier, |
| 83 | sourceCount: plan.sourcePacks.length, |
| 84 | }); |
| 85 | return { |
| 86 | action: "preview", |
| 87 | status: "ok", |
| 88 | queued, |
| 89 | wantedAt: context.wantedAt, |
| 90 | activeCatalog: context.activeCatalog, |
| 91 | packCatalogVersion: context.packCatalogVersion, |
| 92 | plan, |
| 93 | message: "The active pack catalog has compactable tiers.", |
| 94 | }; |
| 95 | } |
| 96 | |
| 97 | /** |
| 98 | * Record a compaction request and schedule background work. |
| 99 | * |
| 100 | * Only queues work when the active catalog has a bounded automatic plan. |
| 101 | * Clears stale `compactionWantedAt` when no safe plan exists. |
| 102 | */ |
| 103 | export async function requestCompactionState(args: { |
| 104 | ctx: DurableObjectState; |
| 105 | env: Env; |
| 106 | prefix: string; |
| 107 | logger?: Logger; |
| 108 | }): Promise<RequestCompactionResult> { |
| 109 | const context = await loadCompactionContext(args); |
| 110 | |
| 111 | if (context.selection.status === "blocked") { |
| 112 | if (typeof context.wantedAt === "number") { |
| 113 | await context.store.delete("compactionWantedAt"); |
| 114 | } |
| 115 | args.logger?.warn("compaction:request-blocked", { |
| 116 | reason: context.selection.blocked.reason, |
| 117 | sourceTier: context.selection.blocked.sourceTier, |
| 118 | activePackCount: context.selection.blocked.activePackCount, |
| 119 | maxSourceObjects: context.selection.blocked.maxSourceObjects, |
| 120 | maxSourceBytes: context.selection.blocked.maxSourceBytes, |
| 121 | smallestWindowObjects: context.selection.blocked.smallestWindowObjects, |
| 122 | smallestWindowBytes: context.selection.blocked.smallestWindowBytes, |
| 123 | }); |
| 124 | return { |
| 125 | action: "request", |
| 126 | status: "blocked", |
| 127 | queued: false, |
| 128 | shouldEnqueue: false, |
| 129 | activeCatalog: context.activeCatalog, |
| 130 | packCatalogVersion: context.packCatalogVersion, |
| 131 | blocked: context.selection.blocked, |
| 132 | reason: context.selection.blocked.reason, |
| 133 | message: |
| 134 | "Automatic compaction is blocked because no safe source-pack window fits the bounded maintenance budget.", |
| 135 | }; |
| 136 | } |
| 137 | |
| 138 | if (context.selection.status === "no_work") { |
| 139 | if (typeof context.wantedAt === "number") { |
| 140 | await context.store.delete("compactionWantedAt"); |
| 141 | } |
| 142 | args.logger?.info("compaction:request-no-work", {}); |
| 143 | return { |
| 144 | action: "request", |
| 145 | status: "no_work", |
| 146 | queued: false, |
| 147 | shouldEnqueue: false, |
| 148 | activeCatalog: context.activeCatalog, |
| 149 | packCatalogVersion: context.packCatalogVersion, |
| 150 | reason: "below-threshold", |
| 151 | message: "The active pack catalog is already within the compaction policy.", |
| 152 | }; |
| 153 | } |
| 154 | |
| 155 | const plan = context.selection.plan; |
| 156 | const wantedAt = Date.now(); |
| 157 | await context.store.put("compactionWantedAt", wantedAt); |
| 158 | await scheduleCompactionWake(args.ctx, args.env); |
| 159 | |
| 160 | args.logger?.info("compaction:request", { |
| 161 | wantedAt, |
| 162 | sourceTier: plan.sourceTier, |
| 163 | targetTier: plan.targetTier, |
| 164 | sourceCount: plan.sourcePacks.length, |
| 165 | }); |
| 166 | return { |
| 167 | action: "request", |
| 168 | status: "queued", |
| 169 | queued: true, |
| 170 | shouldEnqueue: true, |
| 171 | wantedAt, |
| 172 | activeCatalog: context.activeCatalog, |
| 173 | packCatalogVersion: context.packCatalogVersion, |
| 174 | plan, |
| 175 | message: "Recorded a compaction request for this repository and queued background work.", |
| 176 | }; |
| 177 | } |
| 178 | |
| 179 | /** Clear any recorded compaction request without affecting active leases. */ |
| 180 | export async function clearCompactionRequestState(args: { |
| 181 | ctx: DurableObjectState; |
| 182 | logger?: Logger; |
| 183 | }): Promise<ClearCompactionRequestResult> { |
| 184 | const store = asTypedStorage<RepoStateSchema>(args.ctx.storage); |
| 185 | const hadQueuedWork = typeof (await store.get("compactionWantedAt")) === "number"; |
| 186 | await store.delete("compactionWantedAt"); |
| 187 | args.logger?.info("compaction:clear", { |
| 188 | hadQueuedWork, |
| 189 | }); |
| 190 | return { |
| 191 | action: "cleared", |
| 192 | cleared: hadQueuedWork, |
| 193 | message: hadQueuedWork |
| 194 | ? "Cleared the recorded compaction request." |
| 195 | : "No recorded compaction request was present.", |
| 196 | }; |
| 197 | } |