File
Blob: src/worker/tasks/routeCacheSync.ts
| 1 | import { type RepoQueueMessageHandle, type RouteCacheSyncMessage } from "./types"; |
| 2 | |
| 3 | import { findNamespaceById, findRepositoryById } from "@/worker/db/d1/dal"; |
| 4 | import { |
| 5 | deleteRouteCacheRecord, |
| 6 | putRouteCacheRecord, |
| 7 | routeCacheKey, |
| 8 | type RouteCacheRecord, |
| 9 | } from "@/worker/repositories/routeCache"; |
| 10 | import { createQueueTaskContext, retryQueueMessage } from "./context"; |
| 11 | |
| 12 | // State-converging route-cache repair. |
| 13 | // |
| 14 | // The consumer reads D1 by `repositoryId` at execution time and reconciles |
| 15 | // ROUTES KV against the canonical row, regardless of what the message body |
| 16 | // says. That makes delivery order, replays, and stale captured slugs all |
| 17 | // safe: the queue body is read-only metadata, and D1 is the source of truth. |
| 18 | // |
| 19 | // Captured slugs (`namespaceSlug`, `repoSlug`) only matter when the row's |
| 20 | // canonical (namespace, slug) has shifted - either because the row is gone, |
| 21 | // or because a future rename will land. In that case we delete the captured |
| 22 | // key in addition to the canonical action. |
| 23 | |
| 24 | const SYNC_RETRY_DELAY_SECONDS = 30; |
| 25 | |
| 26 | export async function handleRouteCacheSyncMessage( |
| 27 | message: Omit<RepoQueueMessageHandle<RouteCacheSyncMessage>, "body">, |
| 28 | body: RouteCacheSyncMessage, |
| 29 | env: Env, |
| 30 | ctx: ExecutionContext |
| 31 | ): Promise<void> { |
| 32 | const task = createQueueTaskContext({ |
| 33 | env, |
| 34 | ctx, |
| 35 | repoLabel: body.repositoryId, |
| 36 | operation: "route-cache-sync", |
| 37 | subrequestBudget: 25, |
| 38 | }); |
| 39 | const log = task.logFor({ service: "RouteCacheSync" }); |
| 40 | log.debug("route-sync:start", { |
| 41 | repositoryId: body.repositoryId, |
| 42 | namespaceSlug: body.namespaceSlug, |
| 43 | repoSlug: body.repoSlug, |
| 44 | enqueuedAt: body.enqueuedAt, |
| 45 | }); |
| 46 | |
| 47 | try { |
| 48 | const repository = await findRepositoryById(task.db, body.repositoryId); |
| 49 | |
| 50 | // D1 row missing: the repo was deleted (or never existed). Drop the |
| 51 | // captured key. We have no canonical (namespace, slug) to address, but |
| 52 | // the captured pair is what the request that enqueued this message |
| 53 | // observed, so deleting it is what the operator expects to converge. |
| 54 | if (!repository) { |
| 55 | await deleteRouteCacheRecord(env, body.namespaceSlug, body.repoSlug); |
| 56 | log.info("route-sync:end", { action: "missing-d1-delete-captured" }); |
| 57 | message.ack(); |
| 58 | return; |
| 59 | } |
| 60 | |
| 61 | const namespace = await findNamespaceById(task.db, repository.namespaceId); |
| 62 | if (!namespace) { |
| 63 | // Defensive: a repository row pointing to a missing namespace is |
| 64 | // structurally inconsistent. Drop the captured key and ack so we |
| 65 | // don't retry forever. |
| 66 | await deleteRouteCacheRecord(env, body.namespaceSlug, body.repoSlug); |
| 67 | log.warn("route-sync:end", { |
| 68 | action: "missing-namespace-delete-captured", |
| 69 | repositoryNamespaceId: repository.namespaceId, |
| 70 | }); |
| 71 | message.ack(); |
| 72 | return; |
| 73 | } |
| 74 | |
| 75 | const canonicalNamespaceSlug = namespace.slug; |
| 76 | const canonicalRepoSlug = repository.slug; |
| 77 | const capturedKey = routeCacheKey(body.namespaceSlug, body.repoSlug); |
| 78 | const canonicalKey = routeCacheKey(canonicalNamespaceSlug, canonicalRepoSlug); |
| 79 | const capturedDiffersFromCanonical = capturedKey !== canonicalKey; |
| 80 | |
| 81 | if (repository.visibility === "private") { |
| 82 | // Private rows must never expose a public route candidate. |
| 83 | if (capturedDiffersFromCanonical) { |
| 84 | await deleteRouteCacheRecord(env, body.namespaceSlug, body.repoSlug); |
| 85 | } |
| 86 | await deleteRouteCacheRecord(env, canonicalNamespaceSlug, canonicalRepoSlug); |
| 87 | log.info("route-sync:end", { |
| 88 | action: "private-delete", |
| 89 | canonicalNamespaceSlug, |
| 90 | canonicalRepoSlug, |
| 91 | deletedCaptured: capturedDiffersFromCanonical, |
| 92 | }); |
| 93 | message.ack(); |
| 94 | return; |
| 95 | } |
| 96 | |
| 97 | // visibility === "public" |
| 98 | if (capturedDiffersFromCanonical) { |
| 99 | await deleteRouteCacheRecord(env, body.namespaceSlug, body.repoSlug); |
| 100 | } |
| 101 | const record: RouteCacheRecord = { |
| 102 | repositoryId: repository.id, |
| 103 | namespaceId: repository.namespaceId, |
| 104 | doName: repository.doName, |
| 105 | updatedAt: repository.updatedAt, |
| 106 | }; |
| 107 | await putRouteCacheRecord(env, canonicalNamespaceSlug, canonicalRepoSlug, record); |
| 108 | log.info("route-sync:end", { |
| 109 | action: "public-put", |
| 110 | canonicalNamespaceSlug, |
| 111 | canonicalRepoSlug, |
| 112 | deletedCaptured: capturedDiffersFromCanonical, |
| 113 | }); |
| 114 | message.ack(); |
| 115 | } catch (error) { |
| 116 | log.warn("route-sync:retry", { error: String(error) }); |
| 117 | retryQueueMessage(message, SYNC_RETRY_DELAY_SECONDS); |
| 118 | } |
| 119 | } |