Skip to content
File

Blob: src/worker/durable/project-do/webhooks/deliveries.ts

typescript285 lines
1import { BranchName, CommitSha, DispatchMode, ExecutionRuntime, RunId, UnixTimestampMs } from "@/contracts";
2import {
3 type EnsureRunInput,
4 type AcceptQueuedRunInput,
5 type RecordVerifiedWebhookDeliveryInput,
6 type RecordVerifiedWebhookDeliveryResult,
7 expectTrusted,
8 nullableTrusted,
9} from "@/worker/contracts";
10import * as projectSchema from "@/worker/db/durable/schema/project-do";
11import { getProjectConfigRow } from "../repo";
12import { ensureRunInitializedWithPayload } from "../run-do-sync";
13import { rescheduleAlarmInTransaction } from "../sidecar-state";
14import { ensureProjectState, transitionAcceptQueuedRun } from "../transitions";
15import type { ProjectDoContext, ProjectStore } from "../types";
16import { WEBHOOK_DELIVERY_RETENTION_MS, type ParsedWebhookDeliveryRow, type StoredWebhookReplayResult } from "./types";
17import {
18 getWebhookDeliveryRow,
19 getProjectWebhookRow,
20 parseWebhookDeliveryRow,
21 pruneWebhookDeliveries,
22 updateWebhookDeliveryRow,
23} from "./repo";
24interface 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}
33interface StaleVerificationMutationResult {
34 outcome: RecordVerifiedWebhookDeliveryResult["outcome"];
35 duplicate: false;
36 runId: null;
37 queuedAt: null;
38 executable: null;
39 staleVerification: true;
40 runInitialization: null;
41}
42interface PreparedWebhookDeliveryMutation {
43 existingRow: typeof projectSchema.projectWebhookDeliveries.$inferSelect | undefined;
44 replay: StoredWebhookReplayResult | null;
45}
46const 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};
50const 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});
61const 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});
81const 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};
107const 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};
121const 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};
144const 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});
158const 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};
255export 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};