File
Blob: src/worker/do/repo/repoDO.ts
| 1 | import type { Head } from "./repoState"; |
| 2 | import type { RepoActivity } from "@/worker/common"; |
| 3 | import type { PackCatalogRow } from "./db/schema"; |
| 4 | |
| 5 | import { DurableObject } from "cloudflare:workers"; |
| 6 | |
| 7 | import { doPrefix } from "@/worker/keys"; |
| 8 | import { text, createLogger } from "@/worker/common"; |
| 9 | import { clearRepositoryStorage, removePack, type RemovePackResult } from "./packOperations"; |
| 10 | import { |
| 11 | abortCompactionLease, |
| 12 | abortReceiveLease, |
| 13 | type BeginCompactionResult, |
| 14 | beginCompactionState, |
| 15 | beginReceiveLease, |
| 16 | type ClearCompactionRequestResult, |
| 17 | clearCompactionRequestState, |
| 18 | clearExpiredLeases, |
| 19 | type CommitCompactionResult, |
| 20 | commitCompactionState, |
| 21 | finalizeReceiveState, |
| 22 | type PreviewCompactionResult, |
| 23 | previewCompactionState, |
| 24 | type RequestCompactionResult, |
| 25 | requestCompactionState, |
| 26 | rearmCompactionQueueFromAlarm, |
| 27 | getActivePackCatalogSnapshot, |
| 28 | getRepoActivitySnapshot, |
| 29 | } from "./catalog"; |
| 30 | import { getRefs, setRefs, resolveHead, setHead, getHeadAndRefs } from "./refs"; |
| 31 | import { handleIdleAndMaintenance } from "./maintenance"; |
| 32 | import { |
| 33 | debugState, |
| 34 | debugCheckCommit, |
| 35 | debugCheckOid, |
| 36 | type DebugCommitCheck, |
| 37 | type DebugOidCheck, |
| 38 | type DebugStateSnapshot, |
| 39 | } from "./debug"; |
| 40 | import { migrate } from "drizzle-orm/durable-sqlite/migrator"; |
| 41 | import { getDb } from "./db"; |
| 42 | import migrations from "../../../../drizzle/repo-do/migrations.js"; |
| 43 | import { |
| 44 | ensureAccessAndAlarm, |
| 45 | touchAndMaybeSchedule, |
| 46 | type RepoDOAccessContext, |
| 47 | } from "./repoDO/access"; |
| 48 | import { seedMinimalRepoState } from "./repoDO/seeding"; |
| 49 | |
| 50 | /** |
| 51 | * Repository Durable Object (per-repo authority) |
| 52 | * |
| 53 | * Responsibilities |
| 54 | * - Acts as the strongly consistent source of truth for repository metadata. |
| 55 | * - Stores refs, HEAD, and pack catalog state in DO storage/SQLite. |
| 56 | * - All operations are provided as typed RPC methods on the class. |
| 57 | * |
| 58 | * Read Path (RPC) |
| 59 | * - Correctness reads live in worker-local pack-first helpers under `src/worker/git/object-store/`. |
| 60 | * - There is no public HTTP endpoint for object reads; this keeps internal state access typed |
| 61 | * and easy to audit. |
| 62 | * |
| 63 | * Write Path |
| 64 | * - Streaming receive: the Worker writes staged `.pack` and `.idx` data to R2, then |
| 65 | * commits refs and pack-catalog metadata through typed DO RPCs. |
| 66 | * |
| 67 | * Maintenance & Background Work |
| 68 | * - `alarm()` handles: lease cleanup, compaction queue re-arm, idle cleanup. |
| 69 | * - The DO is the metadata authority; the data plane lives in R2 packs. |
| 70 | */ |
| 71 | export class RepoDurableObject extends DurableObject { |
| 72 | declare env: Env; |
| 73 | private lastAccessMemMs: number | undefined; |
| 74 | |
| 75 | constructor(ctx: DurableObjectState, env: Env) { |
| 76 | super(ctx, env); |
| 77 | ctx.blockConcurrencyWhile(async () => { |
| 78 | this.lastAccessMemMs = await ctx.storage.get("lastAccessMs"); |
| 79 | const db = getDb(ctx.storage); |
| 80 | await migrate(db, migrations); |
| 81 | // The constructor also runs before `alarm()`. Do not touch |
| 82 | // `lastAccessMs` here, or an alarm wakeup would make an idle object look |
| 83 | // freshly accessed and would keep cleanup from ever seeing it as idle. |
| 84 | }); |
| 85 | } |
| 86 | |
| 87 | async fetch(request: Request): Promise<Response> { |
| 88 | try { |
| 89 | await this.touchAndMaybeSchedule(); |
| 90 | } catch {} |
| 91 | this.logger.debug("fetch", { path: new URL(request.url).pathname, method: request.method }); |
| 92 | return text("Not found\n", 404); |
| 93 | } |
| 94 | |
| 95 | async alarm(): Promise<void> { |
| 96 | this.logger.debug("alarm:start", {}); |
| 97 | |
| 98 | await clearExpiredLeases(this.ctx, this.logger); |
| 99 | |
| 100 | if ( |
| 101 | await rearmCompactionQueueFromAlarm({ ctx: this.ctx, env: this.env, logger: this.logger }) |
| 102 | ) { |
| 103 | return; |
| 104 | } |
| 105 | |
| 106 | await handleIdleAndMaintenance(this.ctx, this.env, this.logger); |
| 107 | this.logger.debug("alarm:end", {}); |
| 108 | } |
| 109 | |
| 110 | private async touchAndMaybeSchedule(): Promise<void> { |
| 111 | await touchAndMaybeSchedule(this.accessContext()); |
| 112 | } |
| 113 | |
| 114 | private async ensureAccessAndAlarm(): Promise<void> { |
| 115 | await ensureAccessAndAlarm(this.accessContext()); |
| 116 | } |
| 117 | |
| 118 | private accessContext(): RepoDOAccessContext { |
| 119 | return { |
| 120 | ctx: this.ctx, |
| 121 | env: this.env, |
| 122 | logger: this.logger, |
| 123 | getLastAccessMemMs: () => this.lastAccessMemMs, |
| 124 | setLastAccessMemMs: (value) => { |
| 125 | this.lastAccessMemMs = value; |
| 126 | }, |
| 127 | }; |
| 128 | } |
| 129 | |
| 130 | public async listRefs(): Promise<{ name: string; oid: string }[]> { |
| 131 | await this.ensureAccessAndAlarm(); |
| 132 | return await getRefs(this.ctx); |
| 133 | } |
| 134 | |
| 135 | public async setRefs(refs: { name: string; oid: string }[]): Promise<void> { |
| 136 | await this.ensureAccessAndAlarm(); |
| 137 | await setRefs(this.ctx, refs); |
| 138 | } |
| 139 | |
| 140 | public async getHead(): Promise<Head> { |
| 141 | await this.ensureAccessAndAlarm(); |
| 142 | return await resolveHead(this.ctx); |
| 143 | } |
| 144 | |
| 145 | public async setHead(head: Head): Promise<void> { |
| 146 | await this.ensureAccessAndAlarm(); |
| 147 | await setHead(this.ctx, head); |
| 148 | } |
| 149 | |
| 150 | public async getHeadAndRefs(): Promise<{ head: Head; refs: { name: string; oid: string }[] }> { |
| 151 | await this.ensureAccessAndAlarm(); |
| 152 | return await getHeadAndRefs(this.ctx); |
| 153 | } |
| 154 | |
| 155 | public async getActivePackCatalog(): Promise<PackCatalogRow[]> { |
| 156 | await this.ensureAccessAndAlarm(); |
| 157 | return await getActivePackCatalogSnapshot(this.ctx); |
| 158 | } |
| 159 | |
| 160 | public async getRepoActivity(): Promise<RepoActivity | null> { |
| 161 | await this.ensureAccessAndAlarm(); |
| 162 | const snapshot = await getRepoActivitySnapshot(this.ctx); |
| 163 | if (snapshot.state === "idle") return null; |
| 164 | return { |
| 165 | state: snapshot.state, |
| 166 | startedAt: snapshot.lease.createdAt, |
| 167 | expiresAt: snapshot.lease.expiresAt, |
| 168 | }; |
| 169 | } |
| 170 | |
| 171 | public async beginReceive() { |
| 172 | await this.ensureAccessAndAlarm(); |
| 173 | return await beginReceiveLease(this.ctx, this.logger); |
| 174 | } |
| 175 | |
| 176 | public async abortReceive(token: string): Promise<boolean> { |
| 177 | await this.ensureAccessAndAlarm(); |
| 178 | return await abortReceiveLease(this.ctx, token); |
| 179 | } |
| 180 | |
| 181 | public async finalizeReceive(args: { |
| 182 | token: string; |
| 183 | commands: Array<{ oldOid: string; newOid: string; ref: string }>; |
| 184 | stagedPack?: |
| 185 | | { |
| 186 | packKey: string; |
| 187 | packBytes: number; |
| 188 | idxBytes: number; |
| 189 | objectCount: number; |
| 190 | } |
| 191 | | undefined; |
| 192 | }) { |
| 193 | await this.ensureAccessAndAlarm(); |
| 194 | return await finalizeReceiveState({ |
| 195 | ctx: this.ctx, |
| 196 | env: this.env, |
| 197 | token: args.token, |
| 198 | commands: args.commands, |
| 199 | stagedPack: args.stagedPack, |
| 200 | logger: this.logger, |
| 201 | }); |
| 202 | } |
| 203 | |
| 204 | public async beginCompaction(): Promise<BeginCompactionResult> { |
| 205 | await this.ensureAccessAndAlarm(); |
| 206 | return await beginCompactionState({ |
| 207 | ctx: this.ctx, |
| 208 | env: this.env, |
| 209 | prefix: this.prefix(), |
| 210 | logger: this.logger, |
| 211 | }); |
| 212 | } |
| 213 | |
| 214 | public async abortCompaction(token: string): Promise<boolean> { |
| 215 | await this.ensureAccessAndAlarm(); |
| 216 | return await abortCompactionLease(this.ctx, token); |
| 217 | } |
| 218 | |
| 219 | public async commitCompaction(args: { |
| 220 | token: string; |
| 221 | sourcePacks: PackCatalogRow[]; |
| 222 | targetTier: number; |
| 223 | packsetVersion: number; |
| 224 | stagedPack: { |
| 225 | packKey: string; |
| 226 | packBytes: number; |
| 227 | idxBytes: number; |
| 228 | objectCount: number; |
| 229 | }; |
| 230 | }): Promise<CommitCompactionResult> { |
| 231 | await this.ensureAccessAndAlarm(); |
| 232 | return await commitCompactionState({ |
| 233 | ctx: this.ctx, |
| 234 | env: this.env, |
| 235 | token: args.token, |
| 236 | sourcePacks: args.sourcePacks, |
| 237 | targetTier: args.targetTier, |
| 238 | packsetVersion: args.packsetVersion, |
| 239 | stagedPack: args.stagedPack, |
| 240 | logger: this.logger, |
| 241 | }); |
| 242 | } |
| 243 | |
| 244 | public async previewCompaction(): Promise<PreviewCompactionResult> { |
| 245 | await this.ensureAccessAndAlarm(); |
| 246 | return await previewCompactionState({ |
| 247 | ctx: this.ctx, |
| 248 | env: this.env, |
| 249 | prefix: this.prefix(), |
| 250 | logger: this.logger, |
| 251 | }); |
| 252 | } |
| 253 | |
| 254 | public async requestCompaction(): Promise<RequestCompactionResult> { |
| 255 | await this.ensureAccessAndAlarm(); |
| 256 | return await requestCompactionState({ |
| 257 | ctx: this.ctx, |
| 258 | env: this.env, |
| 259 | prefix: this.prefix(), |
| 260 | logger: this.logger, |
| 261 | }); |
| 262 | } |
| 263 | |
| 264 | public async clearCompactionRequest(): Promise<ClearCompactionRequestResult> { |
| 265 | await this.ensureAccessAndAlarm(); |
| 266 | return await clearCompactionRequestState({ |
| 267 | ctx: this.ctx, |
| 268 | logger: this.logger, |
| 269 | }); |
| 270 | } |
| 271 | |
| 272 | public async debugState(): Promise<DebugStateSnapshot> { |
| 273 | await this.ensureAccessAndAlarm(); |
| 274 | return await debugState(this.ctx, this.env); |
| 275 | } |
| 276 | |
| 277 | public async debugCheckCommit(commit: string): Promise<DebugCommitCheck> { |
| 278 | await this.ensureAccessAndAlarm(); |
| 279 | return await debugCheckCommit(this.ctx, this.env, commit); |
| 280 | } |
| 281 | |
| 282 | public async debugCheckOid(oid: string): Promise<DebugOidCheck> { |
| 283 | await this.ensureAccessAndAlarm(); |
| 284 | return await debugCheckOid(this.ctx, this.env, oid); |
| 285 | } |
| 286 | |
| 287 | private prefix() { |
| 288 | return doPrefix(this.ctx.id.toString()); |
| 289 | } |
| 290 | |
| 291 | private get logger() { |
| 292 | return createLogger(this.env.LOG_LEVEL, { |
| 293 | service: "RepoDO", |
| 294 | doId: this.ctx.id.toString(), |
| 295 | }); |
| 296 | } |
| 297 | |
| 298 | public async seedMinimalRepo( |
| 299 | withPack: boolean = true |
| 300 | ): Promise<{ commitOid: string; treeOid: string }> { |
| 301 | await this.ensureAccessAndAlarm(); |
| 302 | return await seedMinimalRepoState({ |
| 303 | ctx: this.ctx, |
| 304 | env: this.env, |
| 305 | prefix: this.prefix(), |
| 306 | withPack, |
| 307 | }); |
| 308 | } |
| 309 | |
| 310 | // DO-only storage clear used by the `repository-delete` queue consumer |
| 311 | // after R2 cleanup. Callers must NOT chain this from a Worker handler that |
| 312 | // also enumerates R2 - that would be the forbidden Worker -> DO -> R2 hop. |
| 313 | public async clearRepositoryStorage(): Promise<{ deletedDO: boolean }> { |
| 314 | return await clearRepositoryStorage(this.ctx, this.env); |
| 315 | } |
| 316 | |
| 317 | public async removePack(packKey: string): Promise<RemovePackResult> { |
| 318 | await this.ensureAccessAndAlarm(); |
| 319 | return await removePack(this.ctx, this.env, packKey); |
| 320 | } |
| 321 | } |