File
Blob: src/worker/tasks/queue.ts
| 1 | import { createLogger } from "@/worker/common"; |
| 2 | |
| 3 | import { handleCompactionDeleteMessage, handleCompactionMessage } from "./compaction"; |
| 4 | import { handlePackRefBackfillMessage } from "./refBackfill"; |
| 5 | import { handleRouteCacheSyncMessage } from "./routeCacheSync"; |
| 6 | import { handleRepositoryDeleteMessage } from "./repositoryDelete"; |
| 7 | import { RepoTaskQueueMessageSchema } from "./types"; |
| 8 | |
| 9 | export type { RepoTaskQueueMessage, RepositoryDeleteMessage, RouteCacheSyncMessage } from "./types"; |
| 10 | |
| 11 | // The queue carries repo lifecycle work as well as maintenance: compaction, |
| 12 | // pack-ref backfill, route-cache repair, and repository deletion. Producers |
| 13 | // use the `REPO_TASKS_QUEUE` binding; the physical queue name remains |
| 14 | // `git-on-cloudflare-repo-maint` for continuity. Schemas live in |
| 15 | // `./types.ts`; this file dispatches. |
| 16 | export async function handleRepoTaskQueue( |
| 17 | batch: MessageBatch<unknown>, |
| 18 | env: Env, |
| 19 | ctx: ExecutionContext |
| 20 | ): Promise<void> { |
| 21 | const log = createLogger(env.LOG_LEVEL, { service: "RepoTaskQueue" }); |
| 22 | for (const message of batch.messages) { |
| 23 | const parsed = RepoTaskQueueMessageSchema.safeParse(message.body); |
| 24 | if (!parsed.success) { |
| 25 | log.warn("queue:malformed-message", { |
| 26 | messageId: message.id, |
| 27 | attempts: message.attempts, |
| 28 | issues: parsed.error.issues.map((i) => ({ path: i.path.join("."), code: i.code })), |
| 29 | }); |
| 30 | message.ack(); |
| 31 | continue; |
| 32 | } |
| 33 | const body = parsed.data; |
| 34 | switch (body.kind) { |
| 35 | case "compaction": |
| 36 | await handleCompactionMessage(message, body, env, ctx); |
| 37 | break; |
| 38 | case "compaction-delete": |
| 39 | await handleCompactionDeleteMessage(message, body, env, ctx); |
| 40 | break; |
| 41 | case "pack-ref-backfill": |
| 42 | await handlePackRefBackfillMessage(message, body, env, ctx); |
| 43 | break; |
| 44 | case "route-cache-sync": |
| 45 | await handleRouteCacheSyncMessage(message, body, env, ctx); |
| 46 | break; |
| 47 | case "repository-delete": |
| 48 | await handleRepositoryDeleteMessage(message, body, env, ctx); |
| 49 | break; |
| 50 | default: { |
| 51 | const _exhaustive: never = body; |
| 52 | void _exhaustive; |
| 53 | } |
| 54 | } |
| 55 | } |
| 56 | } |