Skip to content
File

Blob: src/worker/tasks/repositoryDelete.ts

typescript137 lines
1import { type RepoQueueMessageHandle, type RepositoryDeleteMessage } from "./types";
2 
3import { getRepoStub } from "@/worker/common";
4import { deleteRepositoryById } from "@/worker/db/d1/dal";
5import { deleteRouteCacheRecord } from "@/worker/repositories/routeCache";
6import { countSubrequest } from "@/worker/git/operations/limits";
7import { doPrefix } from "@/worker/keys";
8import { 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 
16const 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.
20const DELETE_SUBREQUEST_BUDGET = 800;
21 
22export 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}