Skip to content
File

Blob: src/worker/tasks/routeCacheSync.ts

typescript120 lines
1import { type RepoQueueMessageHandle, type RouteCacheSyncMessage } from "./types";
2 
3import { findNamespaceById, findRepositoryById } from "@/worker/db/d1/dal";
4import {
5 deleteRouteCacheRecord,
6 putRouteCacheRecord,
7 routeCacheKey,
8 type RouteCacheRecord,
9} from "@/worker/repositories/routeCache";
10import { 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 
24const SYNC_RETRY_DELAY_SECONDS = 30;
25 
26export 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}