File
Blob: src/worker/tasks/repositoryDelete.ts
| 1 | import { type RepoQueueMessageHandle, type RepositoryDeleteMessage } from "./types"; |
| 2 | |
| 3 | import { getRepoStub } from "@/worker/common"; |
| 4 | import { deleteRepositoryById } from "@/worker/db/d1/dal"; |
| 5 | import { deleteRouteCacheRecord } from "@/worker/repositories/routeCache"; |
| 6 | import { countSubrequest } from "@/worker/git/operations/limits"; |
| 7 | import { doPrefix } from "@/worker/keys"; |
| 8 | import { createQueueTaskContext, retryQueueMessage } from "./context"; |
| 9 | |
| 10 | // Repository delete owns D1 row removal, ROUTES KV cleanup, R2 enumeration, |
| 11 | // and DO storage clearing. The request handler only authorizes and |
| 12 | // enqueues; failure of any step retries the message. Keeping D1 deletion |
| 13 | // in the consumer means a queue send failure cannot orphan storage cleanup, |
| 14 | // and replay after partial completion still converges to the empty state. |
| 15 | |
| 16 | const DELETE_RETRY_DELAY_SECONDS = 30; |
| 17 | // Bounded per-message work so a very large repo cannot exhaust the |
| 18 | // 1000-subrequest budget; the message is retried with the same body if we |
| 19 | // run out and the next pass picks up where R2 left off. |
| 20 | const DELETE_SUBREQUEST_BUDGET = 800; |
| 21 | |
| 22 | export async function handleRepositoryDeleteMessage( |
| 23 | message: Omit<RepoQueueMessageHandle<RepositoryDeleteMessage>, "body">, |
| 24 | body: RepositoryDeleteMessage, |
| 25 | env: Env, |
| 26 | ctx: ExecutionContext |
| 27 | ): Promise<void> { |
| 28 | const task = createQueueTaskContext({ |
| 29 | env, |
| 30 | ctx, |
| 31 | repoLabel: body.repositoryId, |
| 32 | operation: "delete", |
| 33 | subrequestBudget: DELETE_SUBREQUEST_BUDGET, |
| 34 | }); |
| 35 | const log = task.logFor({ |
| 36 | service: "RepositoryDelete", |
| 37 | repoId: body.doName, |
| 38 | }); |
| 39 | log.info("repo-delete:start", { |
| 40 | repositoryId: body.repositoryId, |
| 41 | namespaceSlug: body.namespaceSlug, |
| 42 | repoSlug: body.repoSlug, |
| 43 | actor: body.actor, |
| 44 | requestedAt: body.requestedAt, |
| 45 | }); |
| 46 | |
| 47 | const { cacheCtx, limiter } = task; |
| 48 | |
| 49 | // Step 1: D1 row delete. Cascade removes repo-scoped PAT grants; |
| 50 | // namespace-scoped grants are intentionally preserved. |
| 51 | try { |
| 52 | const deleted = await deleteRepositoryById(task.db, body.repositoryId); |
| 53 | if (deleted) { |
| 54 | log.info("repo-delete:d1-deleted"); |
| 55 | } else { |
| 56 | log.info("repo-delete:d1-replay-skip"); |
| 57 | } |
| 58 | } catch (error) { |
| 59 | log.warn("repo-delete:retry", { step: "d1", error: String(error) }); |
| 60 | retryQueueMessage(message, DELETE_RETRY_DELAY_SECONDS); |
| 61 | return; |
| 62 | } |
| 63 | |
| 64 | // Step 2: ROUTES KV cleanup at the captured key. The route-cache-sync |
| 65 | // consumer is the canonical converger but a repo-delete message captured |
| 66 | // these slugs at enqueue time, so dropping them here is safe and avoids |
| 67 | // an extra D1 read. |
| 68 | try { |
| 69 | await deleteRouteCacheRecord(env, body.namespaceSlug, body.repoSlug); |
| 70 | log.info("repo-delete:routes-deleted"); |
| 71 | } catch (error) { |
| 72 | log.warn("repo-delete:retry", { step: "routes", error: String(error) }); |
| 73 | retryQueueMessage(message, DELETE_RETRY_DELAY_SECONDS); |
| 74 | return; |
| 75 | } |
| 76 | |
| 77 | // Step 3: R2 enumerate + delete under do/<doId>/. The DO id is derived |
| 78 | // synchronously from `doName`; no DO subrequest is consumed here. |
| 79 | const doId = env.REPO_DO.idFromName(body.doName).toString(); |
| 80 | const prefix = doPrefix(doId); |
| 81 | let cumulativeR2Deleted = 0; |
| 82 | try { |
| 83 | let cursor: string | undefined; |
| 84 | do { |
| 85 | if (!countSubrequest(cacheCtx, 1)) { |
| 86 | log.warn("repo-delete:budget-exhausted", { step: "r2-list" }); |
| 87 | retryQueueMessage(message, DELETE_RETRY_DELAY_SECONDS); |
| 88 | return; |
| 89 | } |
| 90 | const listing = await limiter.run("r2:repo-delete-list", () => |
| 91 | env.REPO_BUCKET.list({ prefix, cursor }) |
| 92 | ); |
| 93 | const objects = listing.objects ?? []; |
| 94 | if (objects.length > 0) { |
| 95 | if (!countSubrequest(cacheCtx, 1)) { |
| 96 | log.warn("repo-delete:budget-exhausted", { step: "r2-delete" }); |
| 97 | retryQueueMessage(message, DELETE_RETRY_DELAY_SECONDS); |
| 98 | return; |
| 99 | } |
| 100 | const keys = objects.map((object) => object.key); |
| 101 | await limiter.run("r2:repo-delete-batch", () => env.REPO_BUCKET.delete(keys)); |
| 102 | cumulativeR2Deleted += keys.length; |
| 103 | log.info("repo-delete:r2-batch", { |
| 104 | count: keys.length, |
| 105 | cumulative: cumulativeR2Deleted, |
| 106 | }); |
| 107 | } |
| 108 | cursor = listing.truncated ? listing.cursor : undefined; |
| 109 | } while (cursor); |
| 110 | } catch (error) { |
| 111 | log.warn("repo-delete:retry", { step: "r2", error: String(error) }); |
| 112 | retryQueueMessage(message, DELETE_RETRY_DELAY_SECONDS); |
| 113 | return; |
| 114 | } |
| 115 | |
| 116 | // Step 4: clear DO storage + alarm. |
| 117 | try { |
| 118 | if (!countSubrequest(cacheCtx, 1)) { |
| 119 | log.warn("repo-delete:budget-exhausted", { step: "do-clear" }); |
| 120 | retryQueueMessage(message, DELETE_RETRY_DELAY_SECONDS); |
| 121 | return; |
| 122 | } |
| 123 | const stub = getRepoStub(env, body.doName); |
| 124 | await limiter.run("do:repo-delete-clear-storage", () => stub.clearRepositoryStorage()); |
| 125 | log.info("repo-delete:do-cleared"); |
| 126 | } catch (error) { |
| 127 | log.warn("repo-delete:retry", { step: "do", error: String(error) }); |
| 128 | retryQueueMessage(message, DELETE_RETRY_DELAY_SECONDS); |
| 129 | return; |
| 130 | } |
| 131 | |
| 132 | log.info("repo-delete:end", { |
| 133 | cumulativeR2Deleted, |
| 134 | }); |
| 135 | message.ack(); |
| 136 | } |