File
Blob: src/worker/api/public/webhooks/handler.ts
| 1 | import { ProjectId, type WebhookProvider as WebhookProviderValue } from "@/contracts"; |
| 2 | import { |
| 3 | type RecordVerifiedWebhookDeliveryResult, |
| 4 | type WebhookTriggerPayload, |
| 5 | expectTrusted, |
| 6 | } from "@/worker/contracts"; |
| 7 | import { createLogger } from "@/worker/services"; |
| 8 | import { findProjectBySlugs } from "@/worker/db/d1/repositories"; |
| 9 | import { requireWebhookProvider, resolveExpectedWebhookInstanceUrl } from "@/worker/api/webhook-shared"; |
| 10 | import { |
| 11 | type WebhookProviderCatalogEntry, |
| 12 | getWebhookProviderCatalogEntry, |
| 13 | isRepositoryUrlWithinInstance, |
| 14 | } from "@/lib/webhooks"; |
| 15 | import type { AppContext } from "@/worker/hono"; |
| 16 | import { HttpError } from "@/worker/http"; |
| 17 | import { normalizeRepositoryUrl, normalizeWebhookInstanceUrl } from "@/worker/validation"; |
| 18 | |
| 19 | import { assertExpectedWebhookContentType, requirePublicSlug } from "./request"; |
| 20 | import { verifyProviderWebhook } from "./verification"; |
| 21 | |
| 22 | const logger = createLogger("worker.webhooks"); |
| 23 | |
| 24 | const 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 | |
| 38 | const 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 | |
| 56 | const 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 | |
| 79 | const 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 | |
| 98 | export 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 | }; |