File
Blob: src/worker/durable/project-do/webhooks/deliveries.ts
| 1 | import { BranchName, CommitSha, DispatchMode, ExecutionRuntime, RunId, UnixTimestampMs } from "@/contracts"; |
| 2 | import { |
| 3 | type EnsureRunInput, |
| 4 | type AcceptQueuedRunInput, |
| 5 | type RecordVerifiedWebhookDeliveryInput, |
| 6 | type RecordVerifiedWebhookDeliveryResult, |
| 7 | expectTrusted, |
| 8 | nullableTrusted, |
| 9 | } from "@/worker/contracts"; |
| 10 | import * as projectSchema from "@/worker/db/durable/schema/project-do"; |
| 11 | import { getProjectConfigRow } from "../repo"; |
| 12 | import { ensureRunInitializedWithPayload } from "../run-do-sync"; |
| 13 | import { rescheduleAlarmInTransaction } from "../sidecar-state"; |
| 14 | import { ensureProjectState, transitionAcceptQueuedRun } from "../transitions"; |
| 15 | import type { ProjectDoContext, ProjectStore } from "../types"; |
| 16 | import { WEBHOOK_DELIVERY_RETENTION_MS, type ParsedWebhookDeliveryRow, type StoredWebhookReplayResult } from "./types"; |
| 17 | import { |
| 18 | getWebhookDeliveryRow, |
| 19 | getProjectWebhookRow, |
| 20 | parseWebhookDeliveryRow, |
| 21 | pruneWebhookDeliveries, |
| 22 | updateWebhookDeliveryRow, |
| 23 | } from "./repo"; |
| 24 | interface AcceptedMutationResult { |
| 25 | outcome: "accepted" | "queue_full"; |
| 26 | duplicate: boolean; |
| 27 | runId: RecordVerifiedWebhookDeliveryResult["runId"]; |
| 28 | queuedAt: RecordVerifiedWebhookDeliveryResult["queuedAt"]; |
| 29 | executable: RecordVerifiedWebhookDeliveryResult["executable"]; |
| 30 | staleVerification: false; |
| 31 | runInitialization: EnsureRunInput | null; |
| 32 | } |
| 33 | interface StaleVerificationMutationResult { |
| 34 | outcome: RecordVerifiedWebhookDeliveryResult["outcome"]; |
| 35 | duplicate: false; |
| 36 | runId: null; |
| 37 | queuedAt: null; |
| 38 | executable: null; |
| 39 | staleVerification: true; |
| 40 | runInitialization: null; |
| 41 | } |
| 42 | interface PreparedWebhookDeliveryMutation { |
| 43 | existingRow: typeof projectSchema.projectWebhookDeliveries.$inferSelect | undefined; |
| 44 | replay: StoredWebhookReplayResult | null; |
| 45 | } |
| 46 | const isWebhookVerificationStillCurrent = (tx: ProjectStore, input: RecordVerifiedWebhookDeliveryInput): boolean => { |
| 47 | const webhookRow = getProjectWebhookRow(tx, input.projectId, input.payload.provider); |
| 48 | return !!webhookRow && webhookRow.enabled !== 0 && webhookRow.updatedAt === input.verifiedWebhookUpdatedAt; |
| 49 | }; |
| 50 | const toStaleVerificationResult = ( |
| 51 | outcome: RecordVerifiedWebhookDeliveryResult["outcome"], |
| 52 | ): StaleVerificationMutationResult => ({ |
| 53 | outcome, |
| 54 | duplicate: false, |
| 55 | runId: null, |
| 56 | queuedAt: null, |
| 57 | executable: null, |
| 58 | staleVerification: true, |
| 59 | runInitialization: null, |
| 60 | }); |
| 61 | const toReplayResult = (row: ParsedWebhookDeliveryRow): StoredWebhookReplayResult => ({ |
| 62 | outcome: row.outcome, |
| 63 | duplicate: true, |
| 64 | runId: row.runId === null ? null : expectTrusted(RunId, row.runId, "RunId"), |
| 65 | queuedAt: null, |
| 66 | executable: null, |
| 67 | staleVerification: false, |
| 68 | runInitialization: null, |
| 69 | payload: { |
| 70 | provider: row.provider, |
| 71 | deliveryId: row.deliveryId, |
| 72 | eventKind: row.eventKind, |
| 73 | eventName: row.eventName, |
| 74 | repoUrl: row.repoUrl, |
| 75 | ref: row.ref, |
| 76 | branch: nullableTrusted(BranchName, row.branch, "BranchName"), |
| 77 | commitSha: nullableTrusted(CommitSha, row.commitSha, "CommitSha"), |
| 78 | beforeSha: nullableTrusted(CommitSha, row.beforeSha, "CommitSha"), |
| 79 | }, |
| 80 | }); |
| 81 | const insertWebhookDeliveryRow = ( |
| 82 | tx: ProjectStore, |
| 83 | input: RecordVerifiedWebhookDeliveryInput, |
| 84 | outcome: RecordVerifiedWebhookDeliveryResult["outcome"], |
| 85 | runId: string | null, |
| 86 | receivedAt: number, |
| 87 | ): void => { |
| 88 | tx.insert(projectSchema.projectWebhookDeliveries) |
| 89 | .values({ |
| 90 | id: crypto.randomUUID(), |
| 91 | projectId: input.projectId, |
| 92 | provider: input.payload.provider, |
| 93 | deliveryId: input.payload.deliveryId, |
| 94 | eventKind: input.payload.eventKind, |
| 95 | eventName: input.payload.eventName, |
| 96 | outcome, |
| 97 | repoUrl: input.payload.repoUrl, |
| 98 | ref: input.payload.ref, |
| 99 | branch: input.payload.branch, |
| 100 | commitSha: input.payload.commitSha, |
| 101 | beforeSha: input.payload.beforeSha, |
| 102 | runId, |
| 103 | receivedAt, |
| 104 | }) |
| 105 | .run(); |
| 106 | }; |
| 107 | const persistWebhookDeliveryRow = ( |
| 108 | tx: ProjectStore, |
| 109 | existingRow: typeof projectSchema.projectWebhookDeliveries.$inferSelect | undefined, |
| 110 | input: RecordVerifiedWebhookDeliveryInput, |
| 111 | outcome: RecordVerifiedWebhookDeliveryResult["outcome"], |
| 112 | runId: string | null, |
| 113 | receivedAt: number, |
| 114 | ): void => { |
| 115 | if (existingRow) { |
| 116 | updateWebhookDeliveryRow(tx, existingRow.id, input, outcome, runId, receivedAt); |
| 117 | return; |
| 118 | } |
| 119 | insertWebhookDeliveryRow(tx, input, outcome, runId, receivedAt); |
| 120 | }; |
| 121 | const prepareWebhookDeliveryMutation = ( |
| 122 | context: ProjectDoContext, |
| 123 | tx: ProjectStore, |
| 124 | input: RecordVerifiedWebhookDeliveryInput, |
| 125 | currentTime: number, |
| 126 | ): PreparedWebhookDeliveryMutation | StaleVerificationMutationResult => { |
| 127 | ensureProjectState(context, tx, input.projectId); |
| 128 | if (!isWebhookVerificationStillCurrent(tx, input)) { |
| 129 | return toStaleVerificationResult(input.outcome); |
| 130 | } |
| 131 | pruneWebhookDeliveries(tx, input.projectId, currentTime - WEBHOOK_DELIVERY_RETENTION_MS); |
| 132 | const existingRow = getWebhookDeliveryRow(tx, input.projectId, input.payload.provider, input.payload.deliveryId); |
| 133 | if (!existingRow) { |
| 134 | return { |
| 135 | existingRow: undefined, |
| 136 | replay: null, |
| 137 | }; |
| 138 | } |
| 139 | return { |
| 140 | existingRow, |
| 141 | replay: toReplayResult(parseWebhookDeliveryRow(existingRow)), |
| 142 | }; |
| 143 | }; |
| 144 | const toDuplicateDeliveryResult = ( |
| 145 | replay: StoredWebhookReplayResult, |
| 146 | ): Omit<RecordVerifiedWebhookDeliveryResult, "queuedAt" | "executable" | "staleVerification"> & { |
| 147 | queuedAt: null; |
| 148 | executable: null; |
| 149 | staleVerification: false; |
| 150 | } => ({ |
| 151 | outcome: replay.outcome, |
| 152 | duplicate: true, |
| 153 | runId: replay.runId, |
| 154 | queuedAt: null, |
| 155 | executable: null, |
| 156 | staleVerification: false, |
| 157 | }); |
| 158 | const runAcceptedWebhookMutation = async ( |
| 159 | context: ProjectDoContext, |
| 160 | input: RecordVerifiedWebhookDeliveryInput, |
| 161 | currentTime: number, |
| 162 | ): Promise<RecordVerifiedWebhookDeliveryResult> => { |
| 163 | const branch = input.payload.branch; |
| 164 | if (branch === null) { |
| 165 | throw new Error(`Accepted webhook delivery ${input.payload.deliveryId} is missing a branch.`); |
| 166 | } |
| 167 | const transition = await context.ctx.storage.transaction( |
| 168 | async (txn): Promise<StoredWebhookReplayResult | AcceptedMutationResult | StaleVerificationMutationResult> => { |
| 169 | const prepared = prepareWebhookDeliveryMutation(context, context.db, input, currentTime); |
| 170 | if ("staleVerification" in prepared) { |
| 171 | return prepared; |
| 172 | } |
| 173 | const { existingRow, replay } = prepared; |
| 174 | if (replay) { |
| 175 | // queue_full is the only retryable delivery state. Providers resend the |
| 176 | // same delivery id, and anvil does not retry internally, so a prior |
| 177 | // queue_full row must be reprocessed until it settles to a terminal state. |
| 178 | if (replay.outcome !== "queue_full") { |
| 179 | return replay; |
| 180 | } |
| 181 | } |
| 182 | const projectConfigRow = getProjectConfigRow(context.db, input.projectId); |
| 183 | if (!projectConfigRow) { |
| 184 | throw new Error(`Project config ${input.projectId} is missing during webhook acceptance.`); |
| 185 | } |
| 186 | const acceptInput: AcceptQueuedRunInput = { |
| 187 | projectId: input.projectId, |
| 188 | triggerType: "webhook", |
| 189 | triggeredByUserId: null, |
| 190 | branch, |
| 191 | commitSha: input.payload.commitSha, |
| 192 | repoUrl: projectConfigRow.repoUrl, |
| 193 | configPath: projectConfigRow.configPath, |
| 194 | provider: input.payload.provider, |
| 195 | deliveryId: input.payload.deliveryId, |
| 196 | dispatchMode: expectTrusted(DispatchMode, projectConfigRow.dispatchMode, "DispatchMode"), |
| 197 | executionRuntime: expectTrusted(ExecutionRuntime, projectConfigRow.executionRuntime, "ExecutionRuntime"), |
| 198 | }; |
| 199 | const accepted = transitionAcceptQueuedRun(context, context.db, acceptInput, currentTime); |
| 200 | if (accepted.kind === "rejected") { |
| 201 | persistWebhookDeliveryRow(context.db, existingRow, input, "queue_full", null, currentTime); |
| 202 | return { |
| 203 | outcome: "queue_full" as const, |
| 204 | duplicate: false, |
| 205 | runId: null, |
| 206 | queuedAt: null, |
| 207 | executable: null, |
| 208 | staleVerification: false, |
| 209 | runInitialization: null, |
| 210 | }; |
| 211 | } |
| 212 | persistWebhookDeliveryRow(context.db, existingRow, input, "accepted", accepted.runId, currentTime); |
| 213 | await rescheduleAlarmInTransaction(context, txn, input.projectId); |
| 214 | return { |
| 215 | outcome: "accepted" as const, |
| 216 | duplicate: false, |
| 217 | runId: accepted.runId, |
| 218 | queuedAt: accepted.queuedAt, |
| 219 | executable: accepted.executable, |
| 220 | staleVerification: false, |
| 221 | runInitialization: accepted.runInitialization, |
| 222 | }; |
| 223 | }, |
| 224 | ); |
| 225 | if (transition.duplicate || transition.runInitialization === null) { |
| 226 | return { |
| 227 | outcome: transition.outcome, |
| 228 | duplicate: transition.duplicate, |
| 229 | runId: transition.runId, |
| 230 | queuedAt: transition.queuedAt, |
| 231 | executable: transition.executable, |
| 232 | staleVerification: transition.staleVerification, |
| 233 | }; |
| 234 | } |
| 235 | try { |
| 236 | await ensureRunInitializedWithPayload(context, transition.runInitialization); |
| 237 | } catch (error) { |
| 238 | context.logger.error("webhook_run_do_initialize_failed", { |
| 239 | projectId: input.projectId, |
| 240 | runId: transition.runId, |
| 241 | provider: input.payload.provider, |
| 242 | deliveryId: input.payload.deliveryId, |
| 243 | error: error instanceof Error ? error.message : String(error), |
| 244 | }); |
| 245 | } |
| 246 | return { |
| 247 | outcome: transition.outcome, |
| 248 | duplicate: false, |
| 249 | runId: transition.runId, |
| 250 | queuedAt: transition.queuedAt, |
| 251 | executable: transition.executable, |
| 252 | staleVerification: false, |
| 253 | }; |
| 254 | }; |
| 255 | export const recordVerifiedWebhookDelivery = async ( |
| 256 | context: ProjectDoContext, |
| 257 | input: RecordVerifiedWebhookDeliveryInput, |
| 258 | ): Promise<RecordVerifiedWebhookDeliveryResult> => { |
| 259 | const currentTime = expectTrusted(UnixTimestampMs, Date.now(), "UnixTimestampMs"); |
| 260 | if (input.outcome === "accepted") { |
| 261 | return await runAcceptedWebhookMutation(context, input, currentTime); |
| 262 | } |
| 263 | return context.db.transaction((tx) => { |
| 264 | const prepared = prepareWebhookDeliveryMutation(context, tx, input, currentTime); |
| 265 | if ("staleVerification" in prepared) { |
| 266 | return prepared; |
| 267 | } |
| 268 | if (prepared.replay) { |
| 269 | // Only the accepted path is allowed to retry a prior queue_full row. |
| 270 | // If this resend now classifies as a non-accepted outcome, replay the |
| 271 | // stored queue_full result instead of mutating durable audit history. |
| 272 | return toDuplicateDeliveryResult(prepared.replay); |
| 273 | } |
| 274 | persistWebhookDeliveryRow(tx, prepared.existingRow, input, input.outcome, null, currentTime); |
| 275 | return { |
| 276 | outcome: input.outcome, |
| 277 | duplicate: false, |
| 278 | runId: null, |
| 279 | queuedAt: null, |
| 280 | executable: null, |
| 281 | staleVerification: false, |
| 282 | }; |
| 283 | }); |
| 284 | }; |