Skip to content
File

Blob: src/worker/api/public/webhooks/handler.ts

typescript167 lines
1import { ProjectId, type WebhookProvider as WebhookProviderValue } from "@/contracts";
2import {
3 type RecordVerifiedWebhookDeliveryResult,
4 type WebhookTriggerPayload,
5 expectTrusted,
6} from "@/worker/contracts";
7import { createLogger } from "@/worker/services";
8import { findProjectBySlugs } from "@/worker/db/d1/repositories";
9import { requireWebhookProvider, resolveExpectedWebhookInstanceUrl } from "@/worker/api/webhook-shared";
10import {
11 type WebhookProviderCatalogEntry,
12 getWebhookProviderCatalogEntry,
13 isRepositoryUrlWithinInstance,
14} from "@/lib/webhooks";
15import type { AppContext } from "@/worker/hono";
16import { HttpError } from "@/worker/http";
17import { normalizeRepositoryUrl, normalizeWebhookInstanceUrl } from "@/worker/validation";
18 
19import { assertExpectedWebhookContentType, requirePublicSlug } from "./request";
20import { verifyProviderWebhook } from "./verification";
21 
22const logger = createLogger("worker.webhooks");
23 
24const canonicalizeRepositoryIdentity = (provider: WebhookProviderValue, normalizedRepoUrl: string): string => {
25 if (provider !== "github") {
26 return normalizedRepoUrl;
27 }
28 
29 const url = new URL(normalizedRepoUrl);
30 const normalizedPath = url.pathname
31 .split("/")
32 .map((segment) => segment.toLowerCase())
33 .join("/");
34 
35 return `https://${url.host}${normalizedPath}`;
36};
37 
38const requireRepositoryMatch = (
39 provider: WebhookProviderValue,
40 normalizedPayloadRepoUrl: string,
41 normalizedProjectRepoUrl: string,
42 expectedInstanceUrl: string | null,
43): void => {
44 if (expectedInstanceUrl && !isRepositoryUrlWithinInstance(normalizedProjectRepoUrl, expectedInstanceUrl)) {
45 throw new HttpError(500, "invalid_webhook_config", "Webhook config is invalid.");
46 }
47 
48 if (
49 canonicalizeRepositoryIdentity(provider, normalizedPayloadRepoUrl) !==
50 canonicalizeRepositoryIdentity(provider, normalizedProjectRepoUrl)
51 ) {
52 throw new HttpError(403, "webhook_repo_mismatch", "Webhook repository does not match the configured project.");
53 }
54};
55 
56const classifyOutcome = (
57 payload: WebhookTriggerPayload,
58 defaultBranch: string,
59): RecordVerifiedWebhookDeliveryResult["outcome"] => {
60 if (payload.eventKind === "ping") {
61 return "ignored_ping";
62 }
63 
64 if (payload.eventKind !== "push" || payload.branch === null) {
65 return "ignored_event";
66 }
67 
68 if (payload.branch !== defaultBranch) {
69 return "ignored_branch";
70 }
71 
72 if (payload.commitSha === null) {
73 return "ignored_event";
74 }
75 
76 return "accepted";
77};
78 
79const mapDeliveryResponseStatus = (
80 catalog: WebhookProviderCatalogEntry,
81 result: RecordVerifiedWebhookDeliveryResult,
82): number => {
83 if (result.outcome === "queue_full") {
84 return catalog.queueFullResponseStatus;
85 }
86 
87 if (result.duplicate) {
88 return 200;
89 }
90 
91 if (result.outcome === "accepted") {
92 return 202;
93 }
94 
95 return 200;
96};
97 
98export const handleWebhook = async (c: AppContext): Promise<Response> => {
99 const provider = requireWebhookProvider(c.req.param("provider"));
100 const ownerSlug = requirePublicSlug(c.req.param("ownerSlug"), "ownerSlug");
101 const projectSlug = requirePublicSlug(c.req.param("projectSlug"), "projectSlug");
102 const catalog = getWebhookProviderCatalogEntry(provider);
103 
104 assertExpectedWebhookContentType(catalog, c.req.raw);
105 
106 const project = await findProjectBySlugs(c.get("db"), ownerSlug, projectSlug);
107 if (!project) {
108 throw new HttpError(404, "webhook_not_found", "Webhook was not found.");
109 }
110 
111 const projectId = expectTrusted(ProjectId, project.id, "ProjectId");
112 const projectStub = c.env.PROJECT_DO.getByName(projectId);
113 const ingressState = await projectStub.getProjectWebhookIngressState(projectId, provider);
114 if (!ingressState || !ingressState.webhook.enabled) {
115 throw new HttpError(404, "webhook_not_found", "Webhook was not found.");
116 }
117 
118 const body = new Uint8Array(await c.req.raw.arrayBuffer());
119 const payload = await verifyProviderWebhook(provider, catalog, c.env, c.req.raw, ingressState.webhook, body);
120 const normalizedProjectRepoUrl = normalizeRepositoryUrl(ingressState.repoUrl);
121 const rawExpectedInstanceUrl = resolveExpectedWebhookInstanceUrl(provider, ingressState.webhook.config);
122 const expectedInstanceUrl =
123 rawExpectedInstanceUrl === null ? null : normalizeWebhookInstanceUrl(rawExpectedInstanceUrl);
124 requireRepositoryMatch(provider, payload.repoUrl, normalizedProjectRepoUrl, expectedInstanceUrl);
125 
126 const outcome = classifyOutcome(payload, ingressState.defaultBranch);
127 const result = await projectStub.recordVerifiedWebhookDelivery({
128 projectId,
129 payload,
130 outcome,
131 verifiedWebhookUpdatedAt: ingressState.webhook.updatedAt,
132 });
133 if (result.staleVerification) {
134 throw new HttpError(409, "stale_webhook_verification", "Webhook verification material is stale.");
135 }
136 
137 if (result.outcome === "accepted") {
138 c.executionCtx.waitUntil(
139 projectStub.kickReconciliation().catch((error) => {
140 logger.warn("webhook_reconciliation_kick_failed", {
141 projectId,
142 provider,
143 deliveryId: payload.deliveryId,
144 error: error instanceof Error ? error.message : String(error),
145 });
146 }),
147 );
148 }
149 
150 logger.info("webhook_delivery_processed", {
151 projectId,
152 provider,
153 deliveryId: payload.deliveryId,
154 outcome: result.outcome,
155 duplicate: result.duplicate,
156 runId: result.runId,
157 });
158 
159 const headers = new Headers();
160 if (result.outcome === "queue_full" && catalog.queueFullRetryAfterSeconds !== null) {
161 // Providers that support receiver backpressure should retry this delivery later.
162 headers.set("retry-after", catalog.queueFullRetryAfterSeconds);
163 }
164 
165 return new Response(null, { status: mapDeliveryResponseStatus(catalog, result), headers });
166};